Skip to main content
🚀 Taking AI from prototype to production? Find the architecture, GPU, security and governance gaps before they become incidents. Get a Production AI Readiness Assessment
GPT-NL slide: a responsible large language model built from scratch, 450 billion text tokens plus 150 billion code tokens from opt-in, legally accepted and synthetic data
AI

Datatrove: Build an LLM Data Curation Pipeline, Measured

A datatrove LLM data curation pipeline on a messy test corpus: licence filtering, Trafilatura, language ID, Gopher/C4, MinHash dedup and PII, with counts.

LB
Luca Berton
¡ 12 min read

Datatrove is Hugging Face’s library for processing text at pretraining scale. You chain readers, extractors, filters, formatters, deduplication steps and writers into a pipeline, then run it on a laptop or on a Slurm cluster without changing the steps. In this post I build a complete LLM data curation pipeline with it on a small corpus I made messy on purpose, and I count what every stage removes. That count held the main lesson: with default settings, the web-crawl quality filters threw away nearly 30% of the good book text.

The idea came from the MLOps Community Amsterdam sovereignty special on 29 January 2026. In the GPT-NL talk from TNO, the architecture slide showed the curation stages built on Hugging Face’s datatrove, with Parquet datasets between stages, running on SLURM. The data slide said GPT-NL trains on opt-in data, data legally accepted for training LLMs and non-IP-infringing synthetic data. That’s all I know about their pipeline. Everything below is my own lab, built from the datatrove source, not a description of how GPT-NL did it.

GPT-NL slide: a responsible large language model built from scratch, 450 billion text tokens plus 150 billion code tokens, comparable to Llama 2 7B and GPT-3 175B

The GPT-NL data slide from the talk: opt-in data, data legally accepted for training and non-IP-infringing synthetic data.

Versions. macOS 26.6 on an Apple M1 Pro (8 cores, 16 GB), Python 3.12.13, datatrove 0.10.1 (released 30 September 2026), trafilatura 1.11.0, spaCy 3.8.16, fasttext-numpy2-wheel 0.9.2, tokenizers 0.23.2, pyarrow 25.0.1 and warcio 1.8.1. I checked every class and argument against the installed source and the v0.10.1 README and examples on GitHub. The API changes between minor versions, so pin it.

The test corpus: known problems, known counts

You can’t judge a filter without knowing what it should remove. My generator downloads seven public-domain Project Gutenberg books (four English, one each in Dutch, German and French; 3 MB), strips the Gutenberg header and footer, unwraps hard-wrapped lines so each paragraph is one line, and cuts the text into documents of 250–600 words. Then it injects problems:

InjectedCountShould be caught by
Licence cc-by-nc-4.0 or unknown142licence filter
tdm_opt_out: true (text and data mining opt-out)40licence filter
Exact copies on a mirror domain130exact dedup
Near-duplicates (a reposting line, a few swapped words)80MinHash
Navigation and cookie-banner pages60Gopher, C4
One sentence repeated 15–30 times30Gopher repetition
English text with a lorem ipsum line10C4
English text with a JSON snippet15C4
Good English text on a link-farm domain25URL filter
Dutch, German and French documents200language ID
Fake emails and +31 6 phone numbers160 docsPII formatter
Public IPs (8.8.8.8, 1.1.1.1, 9.9.9.9) / TEST-NET IPs40 / 20 docsPII formatter

That gives 1,264 documents in four JSONL files and 222 in a Parquet file, plus a gzipped WARC file with 150 HTML pages wrapped in a nav bar, cookie banner, sidebar and footer. 120 pages have text that exists only in the WARC; 30 repeat a JSONL or Parquet document. 130 pages sit on archive.example.org, playing a source with a signed agreement, and 20 on news.example.com, which has none. URLs, licence labels and PII are synthetic, and a source field lets me trace where every document ended up.

Install, and what the extras don’t cover

python3.12 -m venv .venv
.venv/bin/pip install "datatrove[processing]==0.10.1" \
  pyarrow warcio faust-cchardet python-magic orjson zstandard \
  spacy aiohttp requests
export HF_HOME=$PWD/.hf TLDEXTRACT_CACHE=$PWD/.tldcache

The [processing] extra brings Trafilatura, fastText, NLTK, tldextract, tokenizers and xxhash. Everything after it on that line fixed an error I hit:

  • spaCy. GopherQualityFilter failed with Please install spacy to use en word tokenizer. Datatrove picks a word tokenizer per language from assets/tokenizer_assignment.csv, and English and Dutch map to SpaCyTokenizer. spaCy is in the [multilingual] extra, which also pulls TensorFlow, so I installed it alone. It uses spacy.blank("en"), so you don’t need to download a spaCy model.
  • aiohttp and requests. Model downloads go through fsspec’s HTTP filesystem, which needs both.
  • WARC reading needs warcio, faust-cchardet and python-magic (the [io] extra has them, plus datasets). python-magic needs the libmagic C library (brew install libmagic, apt install libmagic1). I didn’t want system packages here, so I put a stub magic.py on PYTHONPATH. That’s only safe because every record in my file has a WARC-Identified-Payload-Type header, and the reader only calls libmagic when it’s missing. Don’t do that with real crawls.

Datatrove caches models and block lists under $HF_HOME/assets. Put it on a disk with room: the URL filter’s built-in block list unpacks to a 124 MB domains file on first use.

A pipeline is a list of steps that pass Document objects (text, id, metadata) along. The readers’ default adapter moves unknown fields into metadata, so my license column arrives as doc.metadata["license"] without extra code. Readers can be chained, and any generator function with the signature (data, rank, world_size) works as a step.

Stage 0: WARC pages through Trafilatura

CONSENTED_DOMAINS = {"archive.example.org": "agreement"}

def tag_licence_from_domain(data, rank: int = 0, world_size: int = 1):
    """WARC records carry no licence field: look the source up in a registry."""
    for doc in data:
        host = urlparse(doc.metadata["url"]).hostname
        doc.metadata["license"] = CONSENTED_DOMAINS.get(host, "unknown")
        doc.metadata["tdm_opt_out"] = False
        yield doc

extract = LocalPipelineExecutor(
    pipeline=[
        WarcReader("data/raw/warc"),
        tag_licence_from_domain,
        Trafilatura(favour_precision=True, timeout=10.0),
        JsonlWriter("out/extracted"),
    ],
    tasks=1,
    logging_dir="logs/0_extract",
)

WarcReader keeps HTML response records (and WET conversion records) and sets url and date in the metadata. Trafilatura runs trafilatura.extract() in a sandboxed subprocess with a per-document timeout (default 1 second). All 150 pages were extracted, and none of the nav, cookie, sidebar or footer strings survived. The <h1> title did, which matters later.

A crawl has no licence column, so the licence has to come from somewhere else. Here it’s a registry of sources with an agreement, keyed by domain. In a real project that registry would be your contracts or consent records. Once it’s written into the metadata, licence filtering is just another metadata filter.

Stage 1: licence first, then URL, language and quality

ALLOWED_LICENCES = {"public-domain", "cc0-1.0", "cc-by-4.0", "agreement"}

def licence_ok(doc: Document) -> bool:
    return (doc.metadata.get("license") in ALLOWED_LICENCES
            and not doc.metadata.get("tdm_opt_out", False))

def removed(name):
    return JsonlWriter(f"out/removed/{name}")

filtering = LocalPipelineExecutor(
    pipeline=[
        JsonlReader("data/raw/jsonl"),
        ParquetReader("data/raw/parquet"),
        JsonlReader("out/extracted"),
        LambdaFilter(licence_ok, exclusion_writer=removed("1_licence")),
        URLFilter(extra_domains=["linkfarm.example.net"], exclusion_writer=removed("2_url")),
        language_filter,  # see below
        GopherRepetitionFilter(exclusion_writer=removed("4_gopher_rep")),
        GopherQualityFilter(exclusion_writer=removed("5_gopher_qual")),
        C4QualityFilter(filter_no_terminal_punct=False, exclusion_writer=removed("6_c4")),
        FineWebQualityFilter(exclusion_writer=removed("7_fineweb_qual")),
        ParquetWriter("out/filtered", schema=SCHEMA),
    ],
    tasks=4,
    workers=4,
    logging_dir="logs/1_filter",
    depends=extract,
)

The licence and consent check goes first. It only reads metadata, so it’s the cheapest filter, and a document you may not use shouldn’t cost you language ID or tokenization. It’s also the step you’ll have to justify later, so give it its own exclusion_writer. Every filter can write what it drops to a folder, and those folders are your audit trail. Built-in filters add a filter_reason to each dropped document; a LambdaFilter returns only true or false, so its output has none.

URLFilter reads metadata["url"]. extra_domains adds registered domains or full hostnames to the built-in block lists; my link farm matched as dropped_subdomain. For real crawls, run it before Trafilatura, on the raw HTML, as the repo’s FineWeb example does: it’s cheaper than extraction.

filter_no_terminal_punct=False is also the FineWeb setting. With the default True, C4 deletes every line that doesn’t end in ., ?, !, " or '. Its line rules assume one paragraph per line, which is why my generator unwraps paragraphs.

Language ID with the 938 KB model

LanguageFilter defaults to fastText’s lid.176.bin, a 131 MB download; the glotlid backend fetches a bigger model from the Hugging Face Hub. fastText also publishes lid.176.ftz, the compressed version of the same 176-language model, while its docs call the .bin “faster and slightly more accurate”. Datatrove 0.10.1 has no argument for the model URL, but it’s a class attribute, so a subclass is enough:

class FT176Compressed(FT176LID):
    MODEL_URL = "https://dl.fbaipublicfiles.com/fasttext/supervised-models/lid.176.ftz"
    MODEL_SUBFOLDER = "ft176-ftz"

language_filter = LanguageFilter(languages=["en"], exclusion_writer=removed("3_language"))
language_filter.model = FT176Compressed(["en"])

The filter writes language and language_score to the metadata and keeps documents whose English score is above language_threshold (default 0.65). The small model removed all 180 Dutch, German and French documents that passed the licence check, and no English ones. The fastText models are licensed CC BY-SA 3.0; a licence-first project should track that too.

One Parquet schema for three readers

My first run crashed with Table schema does not match schema used to create file. ParquetWriter infers its schema from the first document, and the WARC documents carry a date that the others don’t. The writer’s schema argument fixes it:

SCHEMA = pa.schema([
    ("text", pa.string()),
    ("id", pa.string()),
    ("metadata", pa.struct([
        ("url", pa.string()), ("license", pa.string()), ("tdm_opt_out", pa.bool_()),
        ("source", pa.string()), ("date", pa.string()),
        ("language", pa.string()), ("language_score", pa.float64()), ("token_count", pa.int64()),
    ])),
])

Missing fields become null and unknown ones are dropped, including the absolute file_path that JsonlReader adds by default.

Stage-by-stage counts

Every number comes from the stats.json that each executor writes to its logging_dir. The “tuned” columns are a second run with two settings changed, explained below.

StageDefaults: removedDefaults: leftTuned: removedTuned: left
Input (JSONL + Parquet + WARC)1,6361,636
Licence and consent2021,4342021,434
URL filter251,409251,409
Language ID (English)1801,2291801,229
Gopher repetition661,163661,163
Gopher quality231932171,146
C4 quality21911301,116
FineWeb quality107804111,105
Exact dedup97707133972
MinHash dedup5065765907
GPT-2 tokens, with EOS271,320378,434

The licence step removed 89 unknown documents (69 labelled, plus the 20 pages from the source without an agreement), 73 cc-by-nc-4.0 and 40 opt-outs. Every boilerplate page, repeated sentence, lorem ipsum draft, JSON snippet and link-farm document was gone before deduplication. All eleven executors together took 20–23 seconds with warm caches.

The defaults removed nearly 30% of the good book text

With default settings, GopherQualityFilter dropped 160 English book documents for gopher_below_alpha_threshold. Gopher wants at least 80% of words to contain a letter, and the spaCy tokenizer makes every comma and quotation mark a word. Dialogue-heavy chapters of Alice, Sherlock Holmes and Pride and Prejudice scored between 0.67 and 0.80. The argument is named max_non_alpha_words_ratio, but it acts as a minimum; the source has a TODO to rename it.

FineWebQualityFilter then dropped 78 more for line_punct_ratio: under 12% of lines ending in terminal punctuation. A paragraph of dialogue ends in a closing curly quote (”), which isn’t in the default stop_chars. Both filters were tuned for Common Crawl, not for 19th-century novels. The tuned run changes two arguments:

GopherQualityFilter(max_non_alpha_words_ratio=0.65, exclusion_writer=removed("5_gopher_qual"))
FineWebQualityFilter(stop_chars=tuple(TERMINAL_PUNCTUATION) + ("”", "’"),
                     exclusion_writer=removed("7_fineweb_qual"))

Gopher quality now removes only the 17 boilerplate pages, FineWeb removes 11 documents instead of 107, and the final set grows from 657 to 907 documents (39% more tokens) with the same junk removed. C4 still drops three book documents for curly_bracket, because Gutenberg writes superscripts as M^{rs}. That rule removes code and JSON too, so keep C4 away from code data.

Exact dedup, then MinHash

Deduplication runs in several steps because finding duplicates needs the whole dataset, while signatures and filtering can run per shard. Exact dedup has three:

def text_of(doc: Document) -> str:
    return doc.text

exact_cfg = ExactDedupConfig(content_getter=text_of)

exact_1 = LocalPipelineExecutor(pipeline=[ParquetReader("out/filtered"),
            ExactDedupSignature("out/exact/sigs", config=exact_cfg)],
            tasks=4, logging_dir="logs/2a_exact_sigs", depends=filtering)
exact_2 = LocalPipelineExecutor(pipeline=[ExactFindDedups("out/exact/sigs", "out/exact/dups", config=exact_cfg)],
            tasks=1, logging_dir="logs/2b_exact_find", depends=exact_1)
exact_3 = LocalPipelineExecutor(pipeline=[ParquetReader("out/filtered"),
            ExactDedupFilter("out/exact/dups", config=exact_cfg, exclusion_writer=removed("8_exact_dup")),
            ParquetWriter("out/exact_deduped", schema=SCHEMA)],
            tasks=4, logging_dir="logs/2c_exact_filter", depends=exact_2)

content_getter has no default. The first and last steps must read the same input with the same task count, because duplicates are tracked by file and position. Exact dedup removed 97 documents, 19 of them WARC pages. Those weren’t byte-identical to their JSONL originals after Trafilatura, because of the title line. But C4 dropped that two-word line (min_words_per_line=3), and after that the texts matched. Step order decides what counts as “exact”.

MinHash has four steps, as in the repo’s examples/minhash_deduplication.py:

mh_cfg = MinhashConfig(hash_config=HashConfig(precision=64),
                       num_buckets=14, hashes_per_bucket=8, n_grams=5)

mh_1 = LocalPipelineExecutor(pipeline=[ParquetReader("out/exact_deduped"),
         MinhashDedupSignature("out/minhash/sigs", config=mh_cfg)],
         tasks=4, logging_dir="logs/3a_mh_sigs", depends=exact_3)
mh_2 = LocalPipelineExecutor(pipeline=[MinhashDedupBuckets("out/minhash/sigs", "out/minhash/buckets", config=mh_cfg)],
         tasks=mh_cfg.num_buckets, logging_dir="logs/3b_mh_buckets", depends=mh_1)
mh_3 = LocalPipelineExecutor(pipeline=[MinhashDedupCluster("out/minhash/buckets", "out/minhash/remove_ids", config=mh_cfg)],
         tasks=1, logging_dir="logs/3c_mh_cluster", depends=mh_2)
mh_4 = LocalPipelineExecutor(pipeline=[
         ParquetReader("out/exact_deduped"),
         MinhashDedupFilter("out/minhash/remove_ids", exclusion_writer=removed("9_near_dup")),
         PIIFormatter(),
         TokensCounter(tokenizer_name_or_path="gpt2"),
         DocStats("out/stats"),
         ParquetWriter("out/final", schema=SCHEMA)],
         tasks=4, logging_dir="logs/3d_mh_filter", depends=mh_3)

The bucket step asserts that its task count is divisible by num_buckets, and clustering must run as one task. With 14 buckets of 8 hashes, two documents with 5-gram Jaccard similarity s become a candidate pair with probability 1 − (1 − s⁸)¹⁴:

Jaccard0.50.60.70.750.80.850.9
P(flagged)0.050.210.570.770.920.991.00

My near-duplicates scored 0.85–0.87 against their originals. MinHash formed 50 clusters (65 tuned) and kept one document from each. To check for misses, I computed the exact 5-gram Jaccard similarity for every pair in the final output: the highest was 0.038, and no two texts were identical. The reposts that are still there are the ones whose original had already been filtered out.

PII: what PIIFormatter covers

PIIFormatter replaces email addresses and IP addresses, and nothing else. In the final set, all the synthetic emails were gone, and so were the 28 public resolver IPs that reached it. The 11 TEST-NET addresses (203.0.113.0/24) stayed: with only_remove_public_ips=True, only addresses that Python’s ipaddress reports as is_global are replaced. All 101 phone numbers stayed too. For phone numbers, names or ID numbers you need your own BaseFormatter or a dedicated tool.

The replacements rotate through fixed values (email@example.com, firstname.lastname@example.org and six IPs), so the same text formatted twice came out different in my test. That’s one reason to format PII after deduplication, where the FineWeb example also puts it.

Tokenization and statistics

TokensCounter stores token_count per document: 270,663 GPT-2 tokens for 657 documents. DocStats collects length, whitespace and punctuation ratios as summary, histogram, fqdn and suffix groups, and StatsMerger combines the per-task files. The last stage writes training-ready token files:

tokenize = LocalPipelineExecutor(
    pipeline=[ParquetReader("out/final"),
              DocumentTokenizer("out/tokenized", tokenizer_name_or_path="gpt2", eos_token="<|endoftext|>")],
    tasks=4, logging_dir="logs/4b_tokenize", depends=merge_stats,
)

Each task writes a shuffled .ds file with .index and .metadata files: 271,320 tokens, which is 270,663 plus one end-of-text token per document.

LocalPipelineExecutor: tasks, workers, depends

  • tasks is the number of shards; task rank reads files rank, rank + tasks and so on. With one Parquet file and four tasks, three tasks got no Parquet data. Don’t set more tasks than files.
  • workers is how many tasks run at once (default: all of them). Above 1 it uses multiprocessing with forkserver, so keep the entry point under if __name__ == "__main__":.
  • depends chains executors; run() on the last one ran all eleven in order.
  • skip_completed=True writes a marker per task in logging_dir/completions. My second run printed Not doing anything as all 4 tasks have already been completed and finished in 0.6 seconds. That’s handy after a crash, and a trap after a code change: use a fresh logging_dir. Don’t change tasks when resuming, or the sharding changes.

Scaling out with SlurmPipelineExecutor (not tested)

I have no Slurm cluster in this lab, so this is from the README and source only. The same pipeline list goes into SlurmPipelineExecutor, which writes an sbatch script and submits a job array:

from datatrove.executor import SlurmPipelineExecutor

filtering = SlurmPipelineExecutor(
    job_name="curate_filter",
    pipeline=[...],            # the same steps as above
    tasks=1000,
    workers=200,               # max tasks running at once (-1: no limit)
    time="10:00:00",
    partition="cpu",
    cpus_per_task=1,
    mem_per_cpu_gb=2,
    venv_path="/shared/envs/curation/bin/activate",
    logging_dir="s3://my-bucket/logs/filter",
    slurm_logs_folder="logs/filter/slurm_logs",  # must be local
)

depends= becomes a Slurm job dependency, max_array_size (default 1001) splits big arrays, and randomize_start_duration staggers task starts. Don’t launch it from a compute node. If you already run Slurm for training, see Slurm for GPU clusters and multi-node training on Slurm; a CPU partition for curation next to the GPU partitions is a common setup.

Pitfalls, in the order I hit them

  1. GopherQualityFilter needs spaCy for English, which [processing] doesn’t install.
  2. HTTP downloads need aiohttp and requests; WarcReader needs the libmagic C library.
  3. ParquetWriter infers its schema from the first document, so mixed sources need schema=.
  4. The default language model is 131 MB and the URL block list 124 MB: set HF_HOME.
  5. C4’s defaults delete lines without terminal punctuation, and its curly-bracket rule removes code.
  6. Gopher’s alpha-word ratio and FineWeb’s line-punctuation ratio are harsh on dialogue.
  7. skip_completed makes reruns do nothing until you change logging_dir.
  8. PIIFormatter covers emails and public IPs only.

My take

Datatrove does its job well: the steps are small, the stats are honest, and the same code runs locally and on Slurm. The real work is around it: knowing where each document came from and on what terms, putting that check first, and reading what every filter removed before trusting its defaults. For a licensed, curated corpus like the one GPT-NL described, web-crawl heuristics are a starting point, not a verdict. And when data provenance becomes a compliance question under the EU AI Act, the exclusion folders and stats.json files are your evidence. For the infrastructure side, see digital sovereignty in Europe.

Free 30-min Production AI consultation

Book Now