Chapter 9.5
In this chapter · 8 sections
Data Ingestion, Preprocessing & the Data-Loader Path
The data loader sits in the critical loop of every training step; get its format, sharding, and CPU budget wrong and the GPUs stall, burning depreciation while they wait to be fed.
What you'll decide here
- Which on-disk format you commit the corpus to — sequential shards (WebDataset/MDS) vs columnar (Parquet) vs raw files — because that choice is written into petabytes and is expensive to re-encode mid-program.
- Whether you stream samples directly from object storage or stage them onto a node-local/parallel-FS cache first — the fork that sets your egress bill, your startup latency, and your tolerance for object-store tail latency.
- How many CPU cores and host-DRAM gigabytes per GPU you provision for decode, augmentation, shuffle and prefetch — the GPU:CPU ratio that decides whether the loader keeps the selected rack fed or leaves its accelerators drawing power while they wait.
- Whether preprocessing runs on the CPU host, is offloaded to the GPU (DALI), or is baked once into a pre-decoded format (FFCV/MDS) — and what that does to per-epoch repeatability and storage footprint.
- How you make shuffling and resumption deterministic across an elastic GPU count, so a restart after a failure resumes mid-epoch in seconds rather than re-reading from the top of the corpus.
Every chapter in Part 9 so far has been about storage that sits beside the training loop — the parallel file system that holds the corpus (Chapter 9.2), the NVMe tier and GPUDirect path that move bytes fast (Chapter 9.3), the checkpoint stream that protects against failure (Chapter 9.4). This chapter is about the storage path that sits inside the loop. The data loader supplies the input needed by every optimizer step: the GPU cannot start step N until the batch for step N is decoded, augmented, collated, and resident in device memory. If that pipeline cannot sustain the consumption rate, the accelerator stalls and burns depreciation while it waits.
The loader path converts almost every wrong choice into GPU idle time, and idle GPUs convert into wasted megawatts and stranded depreciation. We trace the runtime pipeline stage by stage through the format fork (sequential vs columnar vs raw), the placement fork (stream vs stage vs cache), and the offload fork (CPU vs GPU vs pre-decoded), then close on the two things people get wrong most often: under-provisioning CPU and DRAM relative to GPUs, and treating multimodal pipelines as if they were the text pipeline with pictures bolted on. The offline corpus-construction work — dedup, filtering, tokenization at corpus scale — is a different machine entirely and lives in Chapter 9.9; the legal regime that governs what may enter the corpus is Chapter 10.10.
The runtime pipeline: what happens between storage and the GPU
A training data loader is a producer-consumer pipeline with the GPU as the consumer. Reading left to right, every sample traverses six stages on its way to becoming part of a batch: (1) fetch a shard or object from the storage tier; (2) decode the container and the payload (JPEG/WebP, audio, video frames, or token records); (3) transform/augment (resize, crop, normalize, tokenize, pad); (4) shuffle across a buffer so consecutive batches are decorrelated; (5) collate samples into a tensor batch; and (6) transfer the batch host-to-device over PCIe/NVLink. Stages 1–5 run on the CPU host (unless explicitly offloaded); stage 6 crosses the bus into HBM.
The governing constraint is a throughput inequality. Let the GPUs consume batches at rate C (samples/s) and the loader produce them at rate P. In the ideal steady-state overlapped pipeline, if P >= C and prefetch absorbs the declared jitter, the GPU need not wait for the loader. If P < C, the average data-starved fraction is (1 - P/C); measured per-step waits also include tails beyond that model, and no amount of fabric bandwidth or HBM capacity recovers it — you are bottlenecked on the slowest of stages 1–6. Keeping P >= C as C climbs with each GPU generation — while the corpus grows into the petabytes and per-sample decode cost, especially for images and video, refuses to fall — is the whole engineering problem of the loader path.
This is why MLPerf Storage frames its training results around accelerator utilization (AU) rather than raw bandwidth: real reads must feed an emulated accelerator-consumption schedule while meeting the validity threshold in the MLPerf Storage v2.0 rules for the August 2025 results: 90% AU for 3D U-Net and ResNet-50, 70% for CosmoFlow. The result is a maximum emulated accelerator count under those rules, not a promise for the proposed model’s MFU or training goodput. A storage tier that misses the real loader’s consumption schedule still starves the most expensive asset in the data center. Preserve the benchmark’s workload, version and cache rules, then qualify batch wait, accepted steps and restart behavior with the real loader. → Chapter 9.3 for the measurement boundary; Chapter 9.8 for installed acceptance.
The format fork: sequential shards vs columnar vs raw files
The first irreversible-ish decision is the on-disk format the corpus is encoded into, because it is written across petabytes and re-encoding a multi-PB corpus is a multi-day, multi-thousand-core job in its own right (→ Chapter 9.9). Three families dominate, and they trade the same three things against each other: sequential read efficiency, random-access / column-projection ability, and metadata/small-file pressure on the storage tier.
Sequential shards — WebDataset (POSIX tar) and MosaicML MDS — are the LLM and large-multimodal default. The corpus is packed into shards of typically 50 MB–1 GB (MDS defaults to ~67 MB; 50–100 MB shards work well across modalities per Databricks/Mosaic guidance), each holding thousands of samples in record order. The win is that the loader issues large sequential reads against the storage tier instead of millions of tiny random reads — which is the difference between a parallel file system running near its streaming bandwidth and the same file system collapsing under metadata-server load from a small-file storm (the small-file problem from Chapter 9.2 is exactly what sharding exists to defeat). The cost: random access within a shard is awkward, and you cannot cheaply read one column of a multi-field record without decoding the whole sample.
Columnar — Apache Parquet — is the data-lake-native and tabular/structured default. Parquet stores data by column with embedded compression and statistics, so a loader can project just the fields it needs and push down filters, and it integrates natively with the lakehouse engines that produced the data (→ Chapter 9.6). For text/token corpora and structured multimodal metadata this is excellent; for raw image/video bytes it is less natural, and naive row-group sizing can reintroduce small-read patterns. Parquet is where the corpus lives during curation; many shops convert hot training splits from Parquet into WebDataset/MDS shards for the run itself.
Raw files (one image/clip per object) is the path of least resistance and the most common scaling mistake. A directory of 2 billion JPEGs, as an illustrative corpus, is trivial to produce and punishing to train from when namespace service is undersized: reading each sample can require several metadata operations plus a tiny random read, which saturates the file system's metadata plane long before its bandwidth, and turns object-store request-rate limits into the binding constraint. Raw files are fine for small datasets and prototyping; at corpus scale they can turn a fast storage tier into a slow one when namespace and small-read overhead exceed its service budget. Pack shards when that overhead is the constraint and the access contract tolerates the change.
| Format | Read pattern | Random access | Metadata pressure | Best fit | Main downside |
|---|---|---|---|---|---|
| WebDataset (tar shards) | Large sequential reads of 50 MB–1 GB shards | Weak (scan within shard) | Low — millions of samples, thousands of shards | LLM & large multimodal pre-training on POSIX/parallel FS or object | No column projection; re-sharding the corpus is costly |
| MosaicML MDS | Sequential, ~67 MB default shards; random sample seek supported | Moderate — indexed sample offsets | Low | Streaming from cloud object with deterministic, mid-epoch-resumable shuffle | Framework-coupled tooling; another encode pass |
| Apache Parquet (columnar) | Row-group reads; column projection + predicate pushdown | Good for fields; poor for raw blobs | Low–moderate (row-group dependent) | Text/token corpora, structured + tabular, lakehouse-native curation | Awkward for raw image/video bytes; row-group tuning matters |
| Raw files (1 object/sample) | Millions of tiny random reads | Native (one file = one sample) | Severe — metadata storm, request-rate limited | Small datasets, prototyping, ad-hoc inspection | Collapses storage metadata plane at corpus scale |
Sharding and shuffling: decorrelation without a small-file storm
Sharding solves the read-pattern problem but creates a statistics problem: if you read shards in order and emit samples in shard order, every batch is drawn from a narrow slice of the corpus, gradients are correlated, and training quality suffers. The fix is to shuffle — but you cannot hold a 50 TB corpus in memory to do a global shuffle, so the standard technique is a two-level shuffle: shuffle the order of shards globally (cheap — there are only thousands of them), then maintain a shuffle buffer of K samples in host DRAM from which each batch is drawn at random as new samples stream in. A buffer of a few thousand to tens of thousands of samples is a candidate to test, not a guarantee of global-shuffle quality. Buffer size trades host memory against decorrelation; corpus source and class ordering decide whether that trade also changes training quality. Test sample mixing and quality because a bounded shuffle is not a uniform permutation of the whole corpus.
The choice with the largest downstream cost is determinism under elastic scale. A 10,000-GPU run will be interrupted — Meta's Llama 3 405B saw an unplanned interruption roughly every three hours over a 54-day window on 16,384 H100s (Meta, Llama 3 paper, 2024). When it restarts, you want it to resume mid-epoch, in seconds, having already consumed the samples it processed and not re-reading from the top — and you want the sample order to be identical regardless of whether the restart happens on 16,384 GPUs or 12,288. Streaming loaders built for this (MosaicML StreamingDataset is the canonical example) make the shuffle elastically deterministic: sample order is a function of the saved seed, shuffle algorithm and consumed-sample position, so supported elastic resumption can replay the same ordered samples for debugging a loss spike. Keep global batch size constant and divisible by the restart device count, hold the canonical-node partition fixed, and restore the loader state, corpus manifest, batching policy and library release; exact sample order does not promise bit-identical floating-point training. A loader that lacks this turns every failure into either a full-epoch replay (egress and idle-GPU cost) or a silently non-deterministic data order (irreproducible training). The checkpoint/restart math this rides on is Chapter 9.4.
Prove elastic resumption with an uninterrupted control and an interrupted run from the same immutable dataset. Save the consumed global sample cursor with model state, not the number of samples already fetched into queues. Hold global batch, canonical-node partition, seed, shuffle algorithm and transform policy fixed; restart on a supported device count and compare the next ordered global sample IDs and augmentation seeds. Prefetched but unconsumed records must neither disappear nor count as completed. Reject a restart count that cannot divide the fixed global batch before launching it. The Mosaic Streaming requirements define that support boundary; floating-point reproducibility remains a separate tolerance test.
Data loaders and the CPU bottleneck: the GPU:CPU:DRAM ratio
The failure mode that strands the most compute is invisible until the GPUs are installed. Stages 2–5 of the pipeline — decode, augment, shuffle, collate — run on the CPU host by default, and they do not get faster when you buy faster GPUs. As accelerator throughput climbs generation over generation (H100 to B200 to GB300 to Vera Rubin), the consumption rate C climbs with it, but the host CPU's ability to produce decoded, augmented batches P stays flat unless you provision for it. The result is a pipeline that was perfectly balanced on the previous GPU generation and is CPU-starved on the next. NVIDIA, AWS, and the FFCV authors all describe the same root cause: "these pipelines — currently executed on the CPU — have become a bottleneck, limiting the performance and scalability of training."
The lever is the GPU:CPU-core and GPU:host-DRAM ratio, decided when you spec the server, and it is workload-dependent in a way that catches people out. A pure-text LLM pipeline is cheap per sample — tokens are small, decode is light, augmentation is trivial — so a modest core count per GPU and a few dozen prefetch workers, if profiling confirms that count keeps pace without exhausting host memory, keep P ahead of C. An image or, far worse, a video pipeline is the opposite: media decode cost is codec-, resolution-, frame-count- and pipeline-specific; benchmark CPU seconds and decoded bytes per sample for the selected workload. The same chassis that comfortably feeds eight GPUs on text will starve them on video. Under-provision cores and you cap goodput at whatever the CPU can produce; over-provision and you pay for idle cores — but idle cores are cheap insurance next to idle GPUs.
Host DRAM is the second half of the ratio, and it serves three simultaneous masters: the shuffle buffer, the prefetch queue (you want several batches staged ahead so a slow read never reaches the GPU), and — if you stage rather than stream — the node-local dataset cache. Under-size DRAM and the shuffle buffer shrinks (worse decorrelation), the prefetch depth shrinks (jitter reaches the GPU), or the cache thrashes back to the storage tier. This is one reason a training node can need substantial host DRAM alongside its HBM: the loader’s shuffle and prefetch buffers must fit while page cache, checkpoint staging, optimizer/offload state and the operating system use memory too. Budget and test those reservations together; the DRAM is not all available to the loader.
The node demands 8 × 500 = 4,000 samples/s. Stored reads are 4,000 × 0.40 MB = 1.6 GB/s; decoded delivery is 4,000 × 2.0 MB = 8.0 GB/s. CPU demand is 4,000 × 0.0020 = 8.0 core-seconds/s; ceil(8.0/0.70) selects 12 worker cores, with training and services already reserved outside this budget. A 0.50 s gap needs 2,000 prefetched decoded samples, or 4.0 GB. Compressed shuffle uses 10,000 × 0.40 MB = 4.0 GB, so loader buffers total 4.0 + 4.0 + 1.0 = 9.0 GB.
With process allowance and checkpoint staging, 9.0 + 6.0 + 40 = 55 GB fits the 64 GB allocation. Select compressed shuffle and qualify those 12 cores with the full pipeline. Flip: decoded shuffle needs 20 GB, raising concurrent use to 4.0 + 20 + 1.0 + 6.0 + 40 = 71 GB, which fails. The crossover for the 10,000-sample shuffle is (64 − 4 − 1 − 6 − 40) GB / 10,000 = 1.3 MB/sample. Above it, keep shuffle compressed, reduce the tested buffer depth or add memory; buying more storage bandwidth does not repair the memory deficit. DALI's CPU/mixed/GPU operator model supplies an offload alternative to profile, not a speedup for these assumed rates. Pass the demand and queue trace to Chapter 9.8.
Scope & caveats
Illustrative inputs; qualify project rates and limits before selection.
The offload fork: CPU vs GPU (DALI) vs pre-decoded (FFCV/MDS)
When the CPU cannot keep up, there are exactly three escapes, and they trade CPU cost, GPU cost, and storage footprint against each other.
Scale out the CPU. Add cores, add prefetch workers, add nodes — the brute-force path. It works until you hit the chassis core ceiling or the point where adding CPU sockets to feed GPUs is worse economics than the alternatives. It keeps the data path simple and the format unchanged, which is its real virtue.
Offload decode/augment to the GPU (NVIDIA DALI). DALI can move JPEG/video decode and augmentation onto the GPU (through nvJPEG-based image decoders and NVDEC where the device and codec support them, alongside selected CPU and mixed operators) and runs its own execution engine to overlap I/O with compute, "offloading data preprocessing to the GPU to eliminate the CPU bottleneck." This is the dominant answer for image and video pipelines where decode is the wall. The tradeoff is honest: you are spending GPU cycles and a slice of HBM bandwidth on preprocessing instead of on the model — justified precisely when the CPU alternative would otherwise leave the GPU idle, and a poor trade when the GPU is already compute-bound.
Pre-decode once into a training-optimized format (FFCV, or baking transforms into MDS). FFCV "combines efficient file formats with asynchronous transfers to maximize GPU utilization" — the idea is to pay the decode and resize cost once, offline, and store samples in a layout the loader can stream with near-zero per-sample CPU. This collapses the runtime CPU cost to almost nothing and is the fastest path for fixed-preprocessing workloads (the classic result is large epoch-time reductions on ImageNet-class training). The cost is paid in two coins: a larger or differently-shaped storage footprint (decoded data is bigger than compressed), and a loss of flexibility — augmentations baked at encode time cannot vary per epoch, so heavy randomized augmentation either stays on the runtime path or must be designed into the format.
A DALI operator is not itself proof of GPUDirect Storage. For CPU decode, host staging can be the intended path; for GPU-ready records, select a loader and storage client that issue supported cuFile operations and repeat the process-counter checks in Chapter 9.3. Keep dtype, shape, sample identity and data checksum fixed between direct and staged runs. Select the integration that reduces real batch wait without changing the training input or overcommitting HBM.
| Strategy | Runtime CPU cost | GPU cost | Storage footprint | Per-epoch augmentation flexibility | Best fit |
|---|---|---|---|---|---|
| Scale out CPU + workers | High (the whole pipeline) | None added | Unchanged | Full — anything the CPU can compute | Text/light pipelines; simplest data path |
| GPU offload (DALI) | Depends on selected CPU, mixed and GPU operators | Adds decode/augment to GPU | Unchanged (compressed) | Full — runs on GPU each epoch | Image/video where decode is the wall and GPU has headroom |
| Pre-decoded (FFCV / baked MDS) | Measured decode, transform and collation cost | None added | Larger (decoded) or re-encoded | Limited — fixed transforms baked at encode time | Fixed-preprocessing, decode-bound, throughput-critical runs |
Streaming vs staging vs caching tiers
Independent of format and offload is the placement question: where do the bytes live relative to the GPU when the run starts? There are three postures, and the fork sets your startup latency, your egress bill, and your exposure to storage tail latency.
Stream from object storage. The loader pulls shards directly from S3/GCS/on-prem object as it needs them, holding only the shuffle and prefetch buffers locally (this is what MosaicML StreamingDataset and WebDataset-over-S3 are built for; AWS publishes loader best practices for exactly this — parallel connections, sharding, and prefetch to keep GPUs fed). The win is avoiding a full-corpus staging step; first-shard fetch and decoder startup still determine time to the first batch; the data can be larger than any local disk. The exposure is object-store tail latency and request-rate limits — a single slow GET, if it reaches the GPU, is a stall — and recurring egress cost when compute and storage sit on a billed cross-zone or cross-region route, or the selected service charges for retrieval (the data-gravity economics of Chapter 9.8). Deep prefetch and many parallel connections are what hide the tail.
Stage to a node-local NVMe or parallel-FS cache, then train. Copy the working set onto fast local storage up front and read from there for the run. This eliminates per-step object-store dependence and tail latency, and it amortizes egress over many epochs — the right call for a corpus read many times (Meta's Research SuperCluster famously fronted training with a multi-tens-of-PB flash cache). The cost is the staging delay before the first step and the local capacity ceiling: if the working set exceeds local NVMe, you cannot stage it whole and must fall back to streaming or a hybrid.
Cache hot, stream cold (the hybrid that wins in practice). Most real pipelines stream from object but interpose a node-local or rack-local cache that holds recently-used shards, so the first epoch pays object-store cost and subsequent epochs are served from flash. This is the placement analogue of the storage hierarchy that runs through all of Part 9, and it is why the parallel-FS and NVMe tiers of Chapter 9.2 and Chapter 9.3 exist between object storage and the GPU.
Multimodal pipelines: where the loader stops being uniform
A text pipeline is close to uniform: every sample is a token sequence of bounded size, decode is trivial, and the batch is a clean rectangle. A multimodal pipeline — interleaved image-text, audio, or video — breaks all three assumptions at once, and treating it as "the text pipeline plus pictures" is a common design mistake.
First, per-sample cost variance explodes. A caption is bytes; the image it describes is a megabyte of JPEG that must be decoded and resized; the video clip is hundreds of frames through a hardware decoder. Within one batch, per-sample decode cost varies with codec, resolution and frame count, so the slowest sample in a batch sets the batch latency — and a few heavy video samples can stall an otherwise-fed pipeline. This is the strongest case for GPU-offloaded decode (DALI's hardware NVDEC path) and for length-aware or cost-aware batching that groups similar-cost samples so the straggler problem is bounded.
Second, modalities arrive at different rates and sizes, so the shard often packs heterogeneous records and the collate step must pad and mask across variable-length sequences and variable-resolution tensors — work that is itself CPU-heavy and a frequent hidden bottleneck. Third, alignment and sampling logic (which frames of a video pair with which text span, how to sample within long clips) lives in the loader and must be deterministic for resumption to work. The net engineering consequence: multimodal runs need a higher CPU:GPU ratio, more aggressive GPU offload, and cost-aware batching — and they expose loader weaknesses that a text-only run would never surface. This is also why the offline prep for multimodal corpora (decoding, frame extraction, re-encoding to uniform shards) is so much heavier than for text, and why that work belongs in the offline supercomputer of Chapter 9.9 rather than the runtime path.
Deep dive: how to actually diagnose a starved GPU (and why people blame the wrong subsystem)
The symptom is always the same — GPU utilization sawtooths or sits below target, MFU is disappointing — and the instinct is to blame the fabric or the storage bandwidth. Usually it is the loader, and the diagnosis is a short decision tree. Step 1: is the GPU actually waiting on data? Profile the step. If the GPU shows idle gaps that line up with batch boundaries (a stall at the start of each step that disappears mid-step), data starvation is confirmed; if the GPU is busy but slow, the problem is compute or collectives, not the loader, and you are in Chapter 8.5 territory instead.
Step 2: which stage is the bottleneck? Watch CPU utilization on the data-loader workers. If the CPUs are pinned at 100% while the storage tier is loafing, you are decode/augment-bound — the fix is more cores, GPU offload (DALI), or pre-decoding (FFCV). If the CPUs are idle while the GPU starves, you are I/O-bound — the fix is bigger sequential reads (re-shard), deeper prefetch, more parallel connections to object storage, or a local cache. If neither is saturated but the GPU still stalls, you are jitter-bound: a tail-latency read or a single heavy sample is reaching the GPU because the prefetch buffer is too shallow — deepen it. Step 3: confirm with a synthetic loader. Replace the real dataset with a loader that yields pre-made random tensors at infinite rate; if batch waits disappear at the same tensor shape and step work, the real input path contributed to the stall. Trace its stages to locate the bottleneck; synthetic tensors also remove parsing and augmentation. This three-step routine resolves the overwhelming majority of "the GPUs aren't busy" tickets without touching the fabric or buying more storage — and it is the empirical backstop behind the whole goodput thesis of Chapter 9.1.
Data versioning, lineage, and governance
The loader must implement the program’s declared evidence contract: knowing, and being able to prove, exactly which data went into a given model. Three capabilities make that possible, and they are increasingly built into the format and loader rather than bolted on after.
Versioning pins a corpus to an immutable, content-addressed snapshot so that "model X was trained on dataset version Y" is a precise, reproducible statement — and so that re-running training reads byte-identical inputs. Lineage records the transformation graph from raw source through dedup, filtering, and tokenization (→ Chapter 9.9) to the shards the loader consumed, so any sample in the corpus can be traced back to its origin. Governance enforces what may enter the corpus and proves what did: license and consent status, PII handling, opt-outs, and regional restrictions. At frontier scale, training-data provenance is a first-class liability surface: where copyright, privacy or data-residency obligations apply, answering "what was in the corpus?" with evidence is part of meeting them. Set provenance and retention from the applicable regime and the program’s reproducibility requirements. That regime, and the compliance architecture it demands, is the subject of Chapter 10.10; the loader's job is to make the answer cheap to produce by carrying version and lineage metadata in the shards it streams.
Choose format, worker count and buffering from the measured modality, then select CPU, GPU or offline decoding by the resource it releases. A loader passes when it supplies the required samples through jitter and concurrent checkpoint staging and resumes the declared sample order after failure. Choosing a faster format without those memory and correctness budgets can exchange visible I/O wait for silent data repetition or host-memory exhaustion.
Cite this chapter
Fehn, J. (2026). Data Ingestion, Preprocessing & the Data-Loader Path (Chapter 9.5). The Definitive Guide to AI Data Centers. https://aidatacenterguide.com/part-9-storage-and-data/9-5-data-ingestion-preprocessing-and-the-data-loader-path (accessed 2026-09-29).
@misc{aidc-9-5,
author = {Fehn, Jacob},
title = {Data Ingestion, Preprocessing & the Data-Loader Path (Chapter 9.5)},
howpublished = {The Definitive Guide to AI Data Centers},
year = {2026},
url = {https://aidatacenterguide.com/part-9-storage-and-data/9-5-data-ingestion-preprocessing-and-the-data-loader-path},
note = {Accessed 2026-09-29}
}