Est.

Data Preprocessing Pipeline Design for LLM Training Throughput

Careful stage-by-stage design determines whether GPUs train efficiently or starve for tokens.

Senior Writer · · 11 min read
Cover illustration for “Data Preprocessing Pipeline Design for LLM Training Throughput”
Systems Design · October 2, 2026 · 11 min read · 2,376 words

The preprocessing pipeline that feeds a large language model is an infrastructure problem. It is an infrastructure problem with direct consequences for how many GPU-hours get wasted. The canonical pipeline runs as a fixed sequence, filtering, deduplication, PII removal, tokenization, bin-packing, and format serialization, and each of these stages makes an implicit tradeoff among I/O cost, compute cost, and the rate at which usable tokens reach the GPU at training time. Treating any one of those stages carelessly leaves a GPU sitting partially idle while the pipeline behind it tries to catch up, because preprocessing could not deliver tokens fast enough and training shifted from compute-bound to I/O-bound.

This is why the field has converged on a coarse-to-fine design pattern: cheap, rule-based filters run first across the entire corpus, and expensive model-based scoring runs only on whatever survives. That ordering is a throughput decision. It is a throughput decision, because it concentrates the costliest compute on the stage of the corpus where it actually changes the outcome, and it minimizes the cost paid per usable token that eventually reaches the GPU. A pipeline either saturates its GPUs or it starves them, set by engineering choices made stage by stage, long before training begins.

What each pipeline stage costs in I/O, compute, and memory

Every stage in the pipeline carries its own cost signature, and the mistake that leads to poorly sized infrastructure is treating the whole sequence as a single undifferentiated category of "CPU work". Each stage needs to be reasoned about on its own terms, because the resource it consumes most is rarely the same from one stage to the next.

Web crawl ingestion and language identification are dominated by raw I/O and network throughput, with compute demands that stay minimal. The limiting factor is simply how fast documents can be read off storage and routed to the next stage. The format chosen for serialization at this point, typically Parquet or Arrow, sets a columnar access pattern that shapes, without permanently fixing, how every later stage reads the data. That is an early decision whose effects echo through the rest of the pipeline.

Rule-based filtering, the heuristics, length cuts, and script detection that come next, is cheap on every axis that matters: low compute, low memory, and high parallelism, since documents can be evaluated independently with almost no coordination overhead. Its job is to do that filtering aggressively enough that the expensive stages downstream never have to touch bytes that were never going to survive anyway.

Deduplication costs more than either of the stages before it. Fuzzy deduplication methods such as MinHash LSH require building and querying large fingerprint indices, and those indices need to fit in memory or spill to storage fast enough that the spill doesn't become the pipeline's limiting factor. At billion-document scale, the size of that index turns into a genuine storage design question: whether it lives in DRAM, on NVMe, or across a distributed store sets whether the stage can run in parallel or ends up serializing the entire pipeline behind it.

Model-based quality filtering, run with a classifier or a perplexity scorer, is the most compute-intensive stage of the sequence. Scoring billions of documents, even with a small model, demands real GPU time, and the inference overhead has to be amortized through careful batching or the stage becomes prohibitively slow. Tokenization, by contrast, is moderate on compute and light on memory per document, but naive implementations run sequentially in ways that choke throughput long before compute or memory limits are reached.

Deduplication as the Stage Most Likely to Become a Storage Bottleneck at Scale

Deduplication's fingerprint index, whether built on MinHash buckets or a similar structure, has to support high-throughput random reads and writes at the same time, and that is precisely the workload general-purpose storage handles worst. The cost of that workload does not scale linearly with corpus size, it compounds, and the index that fit comfortably in memory at small scale becomes the component most likely to block the rest of the pipeline once document counts climb into the billions.

The mechanism is straightforward: the index eventually outgrows DRAM, and once it spills to storage, the latency of each individual lookup starts to dominate the stage's total runtime. At that point the pipeline waits on storage, lookup by lookup, while deduplication crawls through a corpus that compute resources could otherwise have processed in a fraction of the time.

The storage tier chosen for that index is a decision with direct architectural consequences, not a matter of preference between roughly equivalent options. All-flash NVMe arrays supply the low-latency random-read throughput this workload needs, while object storage such as S3, however suitable it is for holding the corpus itself, cannot serve the random-access lookup pattern the index demands. Choosing the wrong tier here sets whether deduplication can be spread across parallel workers or whether it collapses into a single index server that every worker has to queue behind.

There is a tempting objection to all of this: sample-deduplicate instead of running full fuzzy deduplication, and accept that some duplicates will slip through. That is not a quality shortcut, it is a throughput tradeoff with a cost that simply moves downstream. Training on duplicate documents wastes GPU cycles on redundant gradient updates, so the cost of deduplication gets paid somewhere regardless, either inside the preprocessing pipeline or later, in wasted training compute that is considerably more expensive per hour than the preprocessing stage would have been.

Tokenization throughput and keeping pace with GPU demand

Tokenization looks cheap when measured per document, which is what makes it dangerous at pipeline scale. Naive implementations serialize I/O and compute together, which caps the rate at which token sequences can move on to the next stage regardless of how fast the tokenizer itself runs on a single document.

Shard granularity sets whether this stage keeps pace: the corpus has to be split across enough independent shards that no single worker's I/O time or CPU time becomes the ceiling for the entire pipeline's throughput. Shard granularity set too coarse or too fine raises cost at a different stage of the pipeline. Shards that are too large finish at uneven times, and the fastest GPUs in the training cluster end up waiting on whichever tokenization shard is slowest to flush. Shards that are too small introduce their own tax, since the metadata management and shuffle I/O needed to coordinate thousands of tiny shards consumes bandwidth that would otherwise be delivering token sequences.

The format tokenization writes out matters just as much as the speed at which it runs. Serializing output into columnar formats like Parquet or Arrow lets the bin-packing stage downstream read token sequences with sequential I/O instead of random access, which is exactly the access pattern NVMe drives and parallel file systems are built to serve at full speed. A tokenization stage that instead emits row-oriented or fragmented output forces the bin-packer into random reads, which cuts the effective throughput of the storage layer and adds latency before the GPU ever sees its first batch. Everything tokenization gets right or wrong on this front becomes the raw material bin-packing has to work with, and bin-packing is where all of it either pays off or gets exposed.

Bin-packing as the stage where accumulated pipeline decisions determine GPU utilization

Bin-packing is the stage that directly sets GPU utilization. A best-fit algorithm fills fixed-length sequence buckets to minimize padding, and padding tokens are pure waste: the GPU still executes the compute for them, but they produce no useful gradient signal whatsoever. Every decision made earlier in the pipeline, how aggressively documents were filtered, how sequence lengths were distributed, how tokenization sharded its output, is visible here in a tightly packed bucket or a bucket full of padding.

Mechanically, individual sequences are delimited by an end-of-text token, and the best-fit packer assigns them greedily into fixed-length buckets to maximize fill rate. How tight that fill rate can get depends entirely on the distribution of sequence lengths arriving from tokenization: a wide distribution of lengths gives the packer room to fit sequences together tightly, while a narrow distribution, many sequences clustered near the same length, leaves systematic gaps that no packing algorithm can fully close. This is the reason upstream filtering decisions are not purely data-quality choices. Which documents survive filtering, and at what length, is also a bin-packing efficiency decision with measurable downstream cost.

The payoff for getting this right is concrete. One production run documented in the literature achieved high throughput and strong model FLOPs utilization, with GPUs kept fully fed so that training stayed GPU-bound rather than I/O-bound. That outcome depends on more than the packing algorithm itself. The I/O side of bin-packing is easy to overlook, but the packer has to read token sequences fast enough to fill buckets without ever stalling the training data loader, and that means the storage layer feeding the packer has to sustain sequential read throughput that matches the GPU's token consumption rate. This is precisely where parallel file systems such as Lustre or GPFS, or high-throughput NVMe capable of delivering tens of gigabytes per second of sequential read, stop being optional infrastructure and become directly load-bearing for training throughput. GPUDirect Storage extends this further by letting GPUs read training data straight from NVMe without routing through the CPU, which removes a CPU bottleneck that would otherwise cap how fast packed sequences can reach the GPU. Fast storage, at this stage, is the mechanism that keeps the data loader full enough that the packing algorithm's efficiency actually translates into GPU utilization.

Storage Tier Design Across the Full Pipeline

No single storage tier serves every stage of this pipeline well, because the access patterns involved are genuinely different from one another. Raw corpus ingestion reads sequentially in large blocks, deduplication index lookups are random and small-block, tokenized shard reads are sequential at medium block sizes, and bin-packed training data delivery demands sustained, latency-sensitive sequential throughput. A flat storage architecture built around one of these patterns ends up under-serving the rest.

The tiered architecture the field has converged on reflects that reality directly. Object storage such as S3 or an equivalent holds the raw and cold-archived corpus; a parallel file system such as Lustre or GPFS serves the active preprocessing stages that need concurrent access from many nodes at once; and fast local NVMe holds the final packed training dataset, for high-throughput sequential delivery to GPUs. Making that architecture work at scale without constant manual intervention requires data orchestration, moving datasets automatically between tiers as they move through pipeline stages, because without it engineers end up managing tier placement by hand and frequently get it wrong. NVMe drives on PCIe Gen 5, especially in multi-drive configurations that multiply aggregate sequential read bandwidth, are what make that final delivery tier capable of matching GPU token consumption rates at training scale.

This storage design space has moved from informal best practice to something formally measured. MLPerf Storage v3.0, released in September 2026, now benchmarks this exact class of workload: it adds an S3 data layer alongside the existing POSIX layer and introduces KV cache and vector database workload tests, which signals that storage benchmarking has moved decisively toward capturing the real needs of AI data systems rather than relying on generic HPC access patterns.

The obvious objection to all of this tiering is that cloud object storage is cheap and scalable, so why not run the entire pipeline on it. Object storage fails specifically on random-access latency. Round-trip times for the small-block random reads that deduplication index lookups require run orders of magnitude above what NVMe delivers, and once GPU-hours are factored into the comparison, object storage's throughput per dollar for sustained sequential reads at training scale is also less favorable than local NVMe. The tiered design is the architecture that matches each stage's access pattern to the storage medium actually built to serve it.

What a Well-Engineered Pipeline Looks Like End to End

The most common mistake teams make is treating preprocessing as a one-time offline job rather than a system that has to be engineered for throughput from the start. Compute and storage get provisioned for a single pipeline run, the bottleneck stage only becomes visible once training is already scheduled, and the team ends up scrambling to fix an infrastructure problem under the worst possible time pressure.

Three principles follow directly from the stage-by-stage analysis above. Profiling has to come before provisioning: each stage's throughput should be measured in isolation against a representative sample of the real corpus before any storage architecture gets committed to, because the bottleneck stage is rarely the one engineers expect going in. Storage has to be matched to access pattern rather than to cost alone, since the stage that is I/O-bound on random reads, deduplication, and the stage that is I/O-bound on sequential reads, bin-packed delivery, need different storage tiers, and a single-tier design will under-serve at least one of them. Sequence length distribution has to be treated as an engineering parameter in its own right, not just a data-quality artifact: the distribution of lengths that survive filtering sets the ceiling on how efficiently bin-packing can run.

Applied consistently, the coarse-to-fine principle also keeps the total compute cost of the pipeline as low as it can be: cheap stages run across the entire corpus, expensive stages run only on whatever survives, and the final packed dataset ends up reflecting the cumulative efficiency of every decision made upstream of it. A pipeline that inverts this ordering, running model-based quality scoring before rule-based filtering, spends GPU inference time scoring documents that a simple heuristic would have discarded in milliseconds.

The figure that ultimately audits all of these decisions is GPU utilization itself. If training comes out I/O-bound rather than compute-bound, the preprocessing pipeline has not done its job, regardless of how clean the underlying data is. Model FLOPs utilization measures the preprocessing pipeline that fed the training run just as much as it measures the training run itself, because every stage covered here, filtering, deduplication, tokenization, bin-packing, and the storage tiers underneath all of them, produces or wastes the tokens that utilization figure reflects, and that is the standard each of them ultimately has to be judged against.

Sources

  1. GitHub - haolpku/Awesome-LLM-Data-Preparation: Data Preparation for Large Language Models — a curated companion to our JCST 2026 survey. Covers Pre-training, Continual Pre-training, and Post-training (SFT/RLHF/RLAIF) across collection, filtering, dedup, generation, evaluation.
  2. Jiaxin Zhang
Filed underSystems Design

More in Systems Design