Schemas, packaging and delivery
Streaming Training Data from Object Storage: MDS and Similar Formats
Quick answer
A streaming dataset format packs samples into fixed-size shards plus an index file, so each training node downloads only the shards it needs from S3, GCS or Azure Blob into a bounded local cache. MosaicML Streaming's MDS format is the most complete option for multi-node text and multimodal training because it adds deterministic, elastic shuffling and mid-epoch resumption. WebDataset tar shards are the simpler alternative when sequential reads are enough [1]. Convert supplier Parquet or JSONL once, pin the dataset version in the index metadata, and stream in-region.
By SourceX Editorial · Updated
Why stream instead of copying the whole corpus
Streaming exists because full local copies of a licensed pre-training corpus are slow to stage, expensive to replicate on every node, and hard to account for. A 20 TB text-and-image corpus copied to 64 nodes is 64 uncontrolled copies; streamed from one access-controlled bucket, it is one governed copy plus caches that are evicted as training proceeds.
That matters for licensed data specifically. A license usually defines which records you received, where they may be stored and who may access them, and a single source-of-truth bucket with read-only roles is far easier to audit than scattered NVMe volumes. Streaming also decouples cluster size from storage planning: a job can start reading within minutes instead of waiting for a multi-hour pre-stage.
The trade-offs are real. You now depend on object-store throughput and request rates, a cold cache costs latency at the start of each epoch, and random access across a corpus stored as large sequential files is either impossible or slow. The format you choose decides how badly each of those bites.
How MDS shards and the index file work
MDS works by serializing samples into binary shard files and describing them in an index that every reader loads first. According to the MosaicML Streaming documentation, the MDSWriter takes original samples and writes them into MDS shards of serialized samples; a "stream" is a collection of those shard files in one remote location. The reader, StreamingDataset, points at a remote URI (for example s3://bucket/corpus/v3/train) and a local cache directory, and fetches shards on demand.
In practice the writer is configured with a column schema, a compression codec, optional per-shard hashes and a shard size limit. Typical column encodings in releases current as of October 2026 include str, int, bytes, json, ndarray, jpeg and png; check the encoding list for the exact version you pin, because it has grown over releases. The index (an index.json at the root of the output directory in current releases) lists every shard with its sample count, sizes, column names and encodings, compression and hashes; verify the exact field names against the official docs for your version before writing tooling that parses it.
Three properties make the index central to delivery hygiene:
- It is the shard manifest. A reader that trusts the index will not notice a missing or swapped shard until it hits it, so reconcile index entries against your checksum manifest on receipt (see dataset manifests and checksums).
- It fixes sample counts. Deterministic shuffling depends on the global sample count, so any re-conversion that changes shard boundaries changes the sample order even if the records are identical.
- It travels with the data. Anything you need to reproduce a run later, such as the dataset version ID, should sit next to it rather than in a wiki.
Deterministic shuffle and mid-epoch resumption
Deterministic shuffling means a given seed produces the same global sample order no matter how many GPUs, nodes or dataloader workers read the data. MosaicML.s Streaming documentation [6] states that sample order is identical regardless of GPU, node or worker count, that a checkpoint trained on 64 GPUs can be debugged on 8 GPUs with full reproducibility, and that resumption mid-epoch is immediate. It also shuffles across all samples assigned to a node, not within a small per-process buffer as many iterable loaders do.
The mechanics behind "elastic" determinism are worth understanding before you size a cluster. Streaming partitions the sample space over a fixed number of canonical nodes (num_canonical_nodes) and then maps physical ranks onto that layout; the shuffle algorithm (shuffle_algo, with block-based variants such as py1e and py1br in current releases), shuffle_seed and shuffle_block_size control the permutation. Hold all four constant across restarts, and set num_canonical_nodes deliberately rather than letting it default from the first job's node count if you expect to scale up later.
Resumption works by saving the dataset's state_dict() with the model checkpoint and restoring it before iteration. Common failure modes:
- Changed shard layout. Re-converting the corpus with a different
size_limitor compression produces a different index, so a resumed run silently trains on a different order and may revisit or skip samples. - Changed batch size or seed. Resumption is defined relative to the global batch; altering it at restart breaks the sample accounting.
- Mixed streams without weights pinned. When several streams (for example licensed support transcripts plus licensed engineering docs) are mixed with
proportionorrepeat, any weight change is a new data mixture and should be versioned as such. - Cache eviction under load. A
cache_limitset too low thrashes shards on large-sample multimodal data; setpredownloadand the cache to cover at least a few batches per worker.
MDS vs WebDataset and other streaming formats
The choice between MDS and WebDataset comes down to whether you need random access with deterministic elastic resumption or only fast sequential reads. WebDataset stores samples in numbered POSIX tar shards, addressed with brace notation such as corpus-{000000..012345}.tar, and groups files that share a basename into one sample [1]. Hugging Face notes that shards are often around 1 GB while full datasets can reach multiple terabytes [2]. Our WebDataset tar shards guide covers when to request that format from a supplier.
Illustrative example: invented to show structure; it does not describe an available dataset.
| Criterion | MDS (MosaicML Streaming) | WebDataset tar shards | Parquet read in place | JSONL (gzip/zstd) |
|---|---|---|---|---|
| Random access to a sample | Yes, via index offsets | No, sequential within a shard | Row-group level | No |
| Global shuffle | Deterministic, node-level, seedable | Shard shuffle plus in-memory buffer | DIY | DIY |
| Mid-epoch resumption | Built in via state_dict() | Approximate unless you track shard and offset | DIY | DIY |
| Elastic node count | Yes, via canonical nodes | Shards must divide sensibly across workers | DIY | DIY |
| Multimodal blobs | Typed columns (jpeg, png, bytes) | Native: one file per modality per sample | Binary columns, awkward for large blobs | Base64, inefficient |
| Tooling outside training | Limited; needs the Streaming library | Standard tar | Excellent (Spark, DuckDB, Arrow) | Excellent |
| Best fit | Multi-node pre-training and SFT with restarts | Large image, audio or video corpora read sequentially | Analytics and filtering before conversion | Supplier interchange, small SFT sets |
Other options exist. Lightning's LitData, TFRecord with tf.data, and Hugging Face datasets in streaming mode each solve part of the problem; evaluate them on the same criteria, especially resumption semantics. Parquet stays valuable upstream: because its metadata sits in a footer that records where every column chunk starts [3], engines can filter and project columns cheaply, which makes it the right place to deduplicate and filter before you write training shards. For columnar interchange more broadly, see Apache Arrow and Arrow-based dataset formats.
Converting supplier Parquet or JSONL into MDS
Treat conversion as a one-way, versioned build step from the supplier's delivered format into your training format, never as an edit of the delivery itself. Suppliers of operational data usually deliver JSONL, Parquet or original files with sidecar metadata; JSON Lines requires UTF-8 without a byte-order mark and a valid JSON value on every line [4], so validate encoding and blank lines before conversion rather than discovering them as writer exceptions at shard 9,000.
A conversion pipeline that holds up under audit looks like this:
- Verify the delivery against the supplier manifest and checksums, byte for byte, and record the delivery ID.
- Keep the delivered files immutable in a raw prefix (for example
s3://licensed-raw/{supplier}/{delivery_id}/) with object versioning or a lock policy, and convert into a separate prefix. - Filter and deduplicate in Parquet or Arrow, logging every rule and the record counts removed. Quality filtering choices are covered in pretraining text quality filtering.
- Carry stable record IDs into an MDS column so any sample can be traced back to a source record, which you will need for deletion requests, license scope questions and contamination checks.
- Write shards with a fixed
size_limit, codec and hash, using one writer per partition and merging indexes if you parallelize. - Stamp version metadata next to the index (see the artifact below) and publish the prefix read-only.
Illustrative example: invented to show structure; it does not describe an available dataset.
{
"dataset_version_id": "support-transcripts-en/2026-09-30/r2",
"license_ref": "LIC-2026-0412 schedule B",
"source_delivery_ids": ["DLV-0091", "DLV-0094"],
"source_format": "jsonl+zstd",
"conversion": {
"tool": "mosaicml-streaming",
"tool_version": "PINNED_VERSION",
"columns": {"record_id": "str", "text": "str", "meta": "json"},
"compression": "zstd",
"hashes": ["sha1"],
"size_limit_bytes": 67108864
},
"filters_applied": ["exact_dedup_sha256", "minhash_0.8", "lang_en>=0.9"],
"record_counts": {"delivered": 4120000, "after_filters": 3870511},
"pii_method_reference": "supplier redaction report v2",
"created_utc": "2026-10-01T14:22:05Z"
}
Whether you store this inside the index file or as a sidecar such as dataset_meta.json depends on whether your tooling tolerates extra index fields; a sidecar in the same prefix is the safer default because library upgrades can rewrite the index. If you also track versions with a dedicated tool, the trade-offs between DVC, lakeFS and table formats are compared in versioning tools for received datasets.
Bucket layout, region and egress
Streaming should run in the same cloud region as the bucket, because cross-region and internet reads are billed while same-region S3-to-EC2 reads generally are not. On AWS, same-Region S3-to-EC2 reads are generally not billed as data transfer, but traffic routed through a NAT gateway incurs processing charges, so add an S3 gateway endpoint to the VPC; confirm against current AWS pricing pages [7] before budgeting.
Every epoch re-reads the corpus unless the cache holds it, so a cross-region mistake multiplies by epoch count and node count. A practical layout:
- One prefix per dataset version (
.../corpus/{dataset_version_id}/{split}/), never overwritten; new versions get new prefixes. - Read-only IAM role per training cluster, scoped to the version prefixes it is licensed to use, with data-event logging on. On Google Cloud, Data Access audit logs for Cloud Storage are off by default [8] and must be enabled.
- Requester Pays only when you are the requester by design. In a Requester Pays bucket the reader pays request and download charges and anonymous access is not allowed [5], which can shift streaming costs onto your training account if a supplier hosts the data that way.
- Shard size tuned to the store. Shards in the tens to low hundreds of megabytes keep GET request counts manageable without making each cache miss slow; measure on your own network.
Cross-account patterns for receiving the data in the first place are covered in delivering licensed datasets to your cloud bucket, and moving the raw delivery across clouds in multi-terabyte network transfer tools.
Streaming checklist for licensed corpora
Use this checklist before the first multi-node run on a newly licensed corpus.
Illustrative example: invented to show structure; it does not describe an available dataset.
| Check | Pass condition |
|---|---|
| Delivery verified | Supplier manifest and checksums reconciled; delivery IDs recorded |
| Raw copy immutable | Raw prefix versioned or locked; training reads only converted prefix |
| Record IDs preserved | Every sample traceable to a supplier record |
| Version stamped | dataset_version_id and license reference beside the index |
| Shuffle pinned | Seed, algorithm, block size and canonical nodes recorded in run config |
| Resumption tested | Kill and resume a short job; compare sample IDs before and after |
| Region matched | Compute and bucket in the same region; gateway endpoint in place |
| Access scoped | Read-only role limited to licensed prefixes; read logging on |
| Cache sized | cache_limit and predownload cover several batches per worker |
| Use scope matched | Pre-training, continued pre-training or SFT use is within the license grant |
The last row is not a formality. Pre-training typically needs a broader grant than fine-tuning, which is covered in the rights grant you need for pre-training.
Where SourceX fits
SourceX sources operational datasets from US companies on request, such as support and sales histories, engineering records and documents, and manages licensing and ongoing purchases; data is not held in stock and a request does not guarantee a match. Every dataset is rights-reviewed and delivered under a license that defines records, uses, term and delivery, with personal details removed or replaced before delivery and the method recorded. Delivery runs through private, access-controlled workflows after an executed agreement and supplier approval, and you then build streaming shards in your own environment. The general delivery process is summarized in how licensed data is delivered, and you can describe the corpus you need on the SourceX buyers page. For the wider set of format guides, start at the delivery formats hub or the AI data hub.
Request a licensed corpus you can stream
If your team needs licensed operational text or multimodal data for pre-training or SFT, describe the records, fields, volume and intended uses rather than naming companies. SourceX looks for US businesses that hold matching data, reviews data and licensing permissions, and agrees pricing and allowed uses in a license before anything is delivered. Start a request on the SourceX buyers page.
Sources
- WebDataset project, "webdataset (GitHub repository)". https://github.com/webdataset/webdataset
- Hugging Face, "WebDataset (Hub documentation)". https://huggingface.co/docs/hub/datasets-webdataset
- The Apache Software Foundation (Apache Parquet), "File Format". https://parquet.apache.org/docs/file-format/
- jsonlines.org, "JSON Lines". https://jsonlines.org/
- Amazon Web Services (Amazon S3 User Guide), "Using Requester Pays general purpose buckets for storage transfers and usage". https://docs.aws.amazon.com/AmazonS3/latest/dev/RequesterPaysBuckets.html
- MosaicML Streaming, "Main Concepts". https://docs.mosaicml.com/projects/streaming/en/stable/getting_started/main_concepts.html
- Amazon Web Services, "Amazon S3 pricing". https://aws.amazon.com/s3/pricing/
- Google Cloud, "Cloud Storage audit logs". https://cloud.google.com/storage/docs/audit-logging
Tell us what your models need
Share scope, volume, language, format, timing and licensing requirements.