Most writing about pretraining data curation is about what to keep: which heuristics, which quality classifiers, which deduplication threshold. This article is about how the job runs once the corpus is measured in billions of documents and terabytes of text, and why the expensive stages of that job have moved onto GPUs. The filtering rules themselves are covered in other articles on this site. Here the subject is the pipeline as a distributed system: which stages are bound by compute, which by shuffle and which by storage, how fuzzy and semantic deduplication are laid out across GPUs, what the memory arithmetic looks like, and where these jobs break.

The running example is a crawl of one billion documents, about 4 TB of extracted text, that has to come out deduplicated, quality filtered and sharded for a training run. Every sizing number is an estimate with its assumptions stated; redo it with your own corpus.

The pipeline as a chain of stages

Raw crawlWARC, dumps, reposExtract and normaliseCPU, per partitionHeuristic filtersCPU, cheap, big cutExact deduphash, group byFuzzy dedupMinHash, LSH, CCSemantic dedupembed, k-meansQuality classifiersGPU inferencePII and decontamCPU, rules, n-gramsMix and shardweights, tokeniseTraining shardsobject storagePersisted stage outputs and side tablespartitions, duplicate ID lists, scoreswriteblue: CPU or IO boundgreen: GPU pays offyellow: shuffle bound
Figure 1. A curation pipeline as a chain of persisted stages. Cheap CPU filters run first so that the GPU stages, fuzzy dedup, semantic dedup and model-based classifiers, see as few documents as possible.

A curation pipeline is a sequence of stages. Each stage reads a set of document partitions and writes either a smaller set of partitions or a side table, such as a list of duplicate IDs or a column of quality scores. Treat it like any other batch data system: every stage persists its output and the next stage reads that output, so a failure in stage six never forces you to recompute stages one to five.

Order the stages by cost per document removed. Heuristic filters are nearly free and often remove a large share of a raw crawl, so they run first; exact dedup precedes fuzzy, and fuzzy precedes semantic. Every document dropped early is one you never embed. For the reasoning behind the filters themselves, see filtering, deduplication and mixture design and training data filtering.

Which stages belong on GPUs

A stage gains from GPUs only if it is limited by arithmetic or in-memory data movement. Stages limited by storage reads or single-threaded Python leave accelerators idle. Profile one partition first.

StageUsually bound byGPU benefitNotes
Extraction, normalisationCPU parsing, storage readsLowScale out on CPU nodes; output Parquet
Heuristic filtersCPU, string scanningLow to moderateVectorised string ops help; run before anything heavy
Exact dedupShuffle of hashesModerateA group-by on a 16-byte hash column; mostly network
Fuzzy dedupHashing compute, then shuffleHighSignature computation is a dense matrix job
Semantic dedupEncoder inference, matrix multipliesHighSame hardware profile as model inference
Quality and safety classifiersModel inference, tokenisationHighWatch for the tokenizer becoming the bottleneck
PII, decontamination, shardingCPU, storage writesLowN-gram matching against benchmark sets is CPU-friendly

A common layout is two pools, CPU workers at both ends and GPUs in the middle, with object storage between them, or one cluster whose scheduler, such as Ray, places stages by resource type.

Fuzzy deduplication on GPUs

Fuzzy deduplication finds documents that are copies with small edits: mirrored pages, reposts with a different footer, forum threads quoted in full. The standard method is MinHash with locality-sensitive hashing, and it maps well onto GPUs.

MinHash. Break each document into overlapping character n-grams, called shingles. Hash every shingle with k independent hash functions and keep the minimum value per function. For two documents, the probability that their minimums match for one function equals the Jaccard similarity of their shingle sets. So the fraction of matching entries in two k-length signatures estimates that similarity, and the signature is a fixed size no matter how long the document is.

LSH banding. Comparing every pair of signatures is quadratic. Instead, split the k = b × r signature into b bands of r values each, and hash each band to a bucket key. Two documents become candidates if any one band matches exactly. The probability of that is P(s) = 1 − (1 − sr)b, an S-curve whose steepest point sits near (1/b)1/r.

NeMo Curator 26.02 documents defaults of 24-character n-grams, 20 bands and 13 hashes per band, so 260 hashes in total. The threshold estimate is (1/20)1/13 ≈ 0.79, which matches the roughly 80 percent similarity that its documentation describes. Working the formula through shows how sharp the cut is:

Jaccard similarity sP(candidate) with b=20, r=13
0.5about 0.002
0.7about 0.18
0.8about 0.68
0.9about 0.997

Buckets to edges to components. Every bucket that holds more than one document produces edges between its members. A common trick is to link each member to one anchor document per bucket rather than emitting every pair, which keeps the edge count linear in bucket size. Connected components over the edge list then group transitive near-duplicates, and one document per component survives. Some pipelines verify candidate pairs with an exact Jaccard computation before accepting an edge. That buys precision at the cost of a second pass over the text.

The signature step is where GPUs earn their place. For a batch of documents it is one dense operation: a (k × shingles) matrix of affine hashes followed by a column-wise minimum. This NumPy sketch has the same shape that a GPU kernel computes per batch:

import zlib
import numpy as np

P = (1 << 31) - 1                      # prime modulus; products fit in int64
K, BANDS, ROWS = 260, 20, 13           # K = BANDS * ROWS
rng = np.random.default_rng(42)        # fixed seed: signatures must be reproducible
a = rng.integers(1, P, K, dtype=np.int64)
b = rng.integers(0, P, K, dtype=np.int64)

def shingles(text, n=24):
    return {zlib.crc32(text[i:i + n].encode()) & P
            for i in range(max(1, len(text) - n + 1))}

def signature(text):
    x = np.fromiter(shingles(text), dtype=np.int64)        # (S,)
    h = (a[:, None] * x[None, :] + b[:, None]) % P          # (K, S)
    return h.min(axis=1)                                    # (K,)

def band_keys(doc_id, sig):
    # one (band, key, doc_id) row per band; a group-by on (band, key) yields buckets
    return [(i, hash(sig[i * ROWS:(i + 1) * ROWS].tobytes()), doc_id)
            for i in range(BANDS)]

NeMo Curator 26.02 documents the whole chain (MinHashStage, LSHStage, BucketsToEdges, ConnectedComponents, then a separate removal workflow) behind one class on a GPU Ray cluster:

from nemo_curator.core.client import RayClient
from nemo_curator.stages.deduplication.fuzzy.workflow import FuzzyDeduplicationWorkflow

ray_client = RayClient()
ray_client.start()

fuzzy_workflow = FuzzyDeduplicationWorkflow(
    input_path="input_data/",
    cache_path="./cache",
    output_path="./results",
    text_field="text",
    char_ngrams=24,
    num_bands=20,
    minhashes_per_band=13
)
fuzzy_workflow.run()

Its documented default is perform_removal=False. The workflow writes duplicate IDs, and removal is a separate step, so you can audit what would be deleted first. Copy that split.

Semantic deduplication

Fuzzy deduplication misses documents that say the same thing in different words: one wire story rewritten by twenty outlets, templated product pages, near-identical tutorials. Semantic deduplication, described in the SemDeDup paper (Abbas et al., 2023), finds them in embedding space. It has three steps, each a GPU-shaped job.

  1. Embed every surviving document with a small encoder, then L2-normalise the vectors. This is batch inference, with the same tuning as any inference job: length bucketing, large batches, bf16.
  2. Cluster the embeddings with k-means so that comparisons stay local. GPU k-means over hundreds of millions of vectors is a sequence of large matrix multiplies against the centroid matrix.
  3. Compare within clusters. Compute pairwise cosine similarity inside each cluster, and drop a document when it is closer than 1 − ε to a document you are keeping.

The third step is quadratic in cluster size, so it has to be tiled. This PyTorch sketch sorts each cluster so that the preferred document comes first, for example by quality score, then marks any document that is too similar to an earlier one:

import torch

def near_duplicates(emb, eps=0.07, tile=8192):
    """emb: (n, d) L2-normalised, sorted best-first. Returns a bool mask of docs to drop."""
    n = emb.shape[0]
    drop = torch.zeros(n, dtype=torch.bool, device=emb.device)
    for start in range(0, n, tile):
        q = emb[start:start + tile]                       # (t, d)
        sim = q @ emb[:start + q.shape[0]].T              # only compare with earlier docs
        rows = torch.arange(start, start + q.shape[0], device=emb.device)[:, None]
        cols = torch.arange(sim.shape[1], device=emb.device)[None, :]
        sim = sim.masked_fill(cols >= rows, -1.0)         # mask self and later docs
        drop[start:start + q.shape[0]] = sim.max(dim=1).values > 1 - eps
    return drop

This greedy version also drops a document whose only near-neighbour was itself dropped. That is usually acceptable. The real cost driver is cluster skew: a cluster of two million boilerplate pages means trillions of comparisons. Cap cluster sizes with a second k-means, or remove boilerplate earlier.

Pick ε by reading pairs, not by guessing. Sample pairs from several similarity bands, read them, and choose the band where most pairs are genuinely redundant. A threshold that is too loose removes legitimately different documents that share a template, such as source files with the same license header.

Classifier filtering throughput

Model-based quality, toxicity and domain classifiers are GPU inference jobs. The useful estimate is the standard 2 × parameters × tokens FLOPs per forward pass. Take a 140M-parameter encoder that reads the first 512 tokens of each document: that is about 1.4 × 1011 FLOPs per document. Over 500 million documents it is about 7 × 1019 FLOPs. If a GPU sustains an assumed 150 TFLOP/s on this workload, the job needs roughly 130 GPU-hours. If utilisation is a third of that, which is common, it needs about 400.

Low utilisation in these jobs is rarely the model. Usually it is one of the following:

  • Tokenisation on one CPU thread. Tokenise in parallel workers ahead of the GPU.
  • Padding waste. Sort or bucket documents by token length before batching, so that a batch of short pages is not padded to 512 tokens.
  • Small batches. Inference without gradients fits far larger batches than training.
  • Reading text from object storage per batch. Prefetch partitions to local disk or memory. See GPU data pipelines for the same stall analysis applied to training input.

A cheap linear classifier on CPUs can remove obvious junk first, so the transformer scores only survivors.

Worked example: sizing a billion-document crawl

Here is the one-billion-document crawl in round numbers. All the reduction ratios are assumptions for illustration, not measurements.

StepDocuments inSizing estimate
Heuristics (assume 40% removed)1.0BCPU only; 4 TB read, about 2.4 TB written
Exact dedup (assume 10% removed)600M16-byte hash + 8-byte ID: about 14 GB shuffled
Fuzzy dedup signatures540M540M × 260 × 4 bytes ≈ 560 GB of signatures
LSH, 5 bands per iteration540Mabout 16 bytes per (key, ID) row: about 43 GB shuffled per iteration, 4 iterations
Semantic dedup embeddingsabout 500M384-dim fp16 = 768 bytes each: about 385 GB
Quality classifierabout 450Mabout 120 to 400 GPU-hours, per the estimate above

The signature table exceeds any node's GPU memory, so signatures must be streamed and written out. LSH is a shuffle problem: a few bands per iteration (the NeMo default is 5) bounds each shuffle at the cost of more passes, so lower it if the network is the limit. When storage throughput limits the GPU stages, GPUDirect Storage is worth evaluating, but measure first.

Failure modes

  • Giant buckets. Empty pages, cookie banners and error pages all hash alike, and one bucket with millions of members exhausts memory at the edges stage. Filter very short and boilerplate documents first, and use the anchor-per-bucket edge scheme.
  • Skewed partitions. One 30 GB input file becomes one task that runs out of memory.
  • Over-deletion. A loose semantic threshold quietly removes distinct code files and legal documents. Audit a sample of removed pairs for each source.
  • Idle GPUs. Classifier jobs at 20 percent utilisation are usually tokenizer-bound or storage-bound, not compute-bound.
  • Stale intermediates. A rerun with new parameters that reads an old cache produces plausible but wrong removal lists. Put the parameters in the cache path.
  • Unstable IDs. IDs derived from file offsets change after re-sharding, so a removal list deletes the wrong documents. Assign IDs once at ingestion; see training data curation for lineage practice.
  • Late decontamination. Benchmark n-gram removal that runs after mixing misses copies that were already upsampled. Run it before mixing.
  • Unexplained removals. Without removal counts per stage and per source, and a version record for each output shard set, nobody can tell why a source shrank by 90 percent.
  • Preemption without checkpoints. A stage that writes only at the end loses hours of work when a node is preempted. Write per-partition outputs.

Trade-offs

GPU versus CPU clusters. GPUs win clearly for signatures, embeddings and classifiers, and add little for extraction and heuristics. Curation also competes with training for the same accelerators. Dedup stages need memory capacity more than peak FLOPs, so older or cheaper GPUs are often the economical choice for them.

Fuzzy versus semantic deduplication. Semantic dedup removes more redundancy, but it also removes more useful diversity, and it costs an encoder pass over the whole corpus. Many teams run fuzzy dedup everywhere and semantic dedup only on sources known to be repetitive.

Which duplicate to keep. Keeping the highest-quality member needs scores before dedup, which reverses the cheap-first ordering; keeping an arbitrary one is cheapest.

What to do next

  1. Profile one partition through every stage and record the time spent in CPU, GPU, network and storage.
  2. Move only the compute-bound stages to GPUs; keep extraction and heuristics on CPU workers.
  3. Assign stable document IDs at ingestion and make every dedup stage emit removal lists.
  4. Compute the LSH S-curve for your band settings and confirm the threshold you actually want.
  5. Size the signature and shuffle volumes with the arithmetic above before reserving hardware.
  6. Cap cluster sizes for semantic dedup and pick its threshold by reading sampled pairs.
  7. Fix tokenizer and padding bottlenecks before buying more GPUs for classifiers.
  8. Put parameters in cache paths, log removal counts per source, and version every output shard set.
Key takeaway: Curation at pretraining scale is a batch data system whose expensive middle, fuzzy dedup, semantic dedup and model-based classifiers, belongs on GPUs while extraction and heuristics stay on CPUs. Order stages cheapest first, persist every stage, carry stable IDs and emit removal lists, size signature and shuffle volumes before reserving hardware, cap cluster and bucket sizes, and check thresholds by reading samples rather than trusting defaults.