Skip to content

Extract, normalize, quality filters (generator stages)

Moduledata.02 · build · Python · Pass 3 · 4 to 6 h
You buildpython/corpus/stage.py: Doc, Stage, compose, extract; python/corpus/filter.py: normalize_text, normalize_unicode, length_filter, identify_lang, lang_filter, gopher_reasons, gopher_rules, dup_ngram_char_frac, repetition_filter, ppl_filter, build_filters
Contractcourse/contracts/py/corpus/stage.pyi · course/contracts/py/corpus/filter.pyi · config: the filters table of corpus-config.schema.json
Testscourse/tests/data.02/ (what they check: section 4); fixtures course/fixtures/data.02/ (200 golden documents)
Needsdata.01 read_raw and the fetch Manifest (or --ref-deps) · reading: lang.08
Used bydata.03 exact dedup consumes your stages’ output and builds on Doc and compose · data.04 near dedup, data.05 PII scrub, and data.06 sharding are stages of the same type, and C1 plugs the L2.1 n-gram model into ppl_filter
MilestoneMS-corpus
Optional depthRae et al., Scaling Language Models: Methods, Analysis & Insights from Training Gopher (2021), appendix A.1 (the quality rules, free); Penedo et al., The FineWeb Datasets (2024), section 3; Unicode Standard Annex #15, Unicode Normalization Forms (free); David Beazley, “Generator Tricks for Systems Programmers” (free)
  • A stage is a function from an iterator of documents to an iterator of documents. Written as a generator, a whole pipeline holds one document at a time, whatever the corpus size (test_stages_are_lazy, test_pipeline_memory_stays_flat).
  • Normalization is NFC (not NFKC), one line-ending convention, no control characters, and it is idempotent: normalizing twice changes nothing (test_normalize_rules_one_by_one, test_normalize_is_idempotent).
  • Language identification by stop-word share is a vote: each stop word votes for every language that lists it, and the winner’s share of the votes is the confidence (test_hand_example_normalize_and_identify, test_identify_lang_hand_labelled).
  • The Gopher rules are seven cheap statistics with published thresholds; each one catches a different kind of junk, and every boundary is inclusive (test_gopher_rules_hand_labelled, test_gopher_boundaries_are_the_papers).
  • The configured pipeline is checked end to end on 200 documents whose fate was decided when they were written, byte for byte (test_golden_200_documents).
Terminal window
ol start data.02 # stubs python/corpus/stage.py and filter.py
ol tests data.02 # read the test catalog first
ol check data.02 # needs data.01 passing (or --ref-deps)
ol diff data.02 # after passing: your code against the reference

data.01 left you raw documents: every line of every source, exactly as published. Train a tokenizer or a language model on that and it learns whatever the sources contain: the same story with Windows line endings and with Unix ones, accented letters encoded two different ways, pages in German when you wanted English, lists of hashtags, keyword spam that repeats one sentence twenty times. At TinyStories scale the damage is a worse loss curve; at web scale, FineWeb measured each of these filters moving downstream benchmark scores. And the corpus will not fit in memory for long: C1’s TinyStories is about 2 GB of text, the pipeline runs inside a durable activity with a memory limit, and a stage that builds a list of all documents is the one that gets the pod killed. This module defines the shape every later stage shares, Stage = Callable[[Iterator[Doc]], Iterator[Doc]], and the first filters built on it.

A Doc is immutable: id (<source_id>:<k>, k the document’s position in its source), source_id, text, and meta (a mapping of facts about the document: its URL, license, and, as stages add them, its language and scores). A stage that changes a document yields a new Doc with dataclasses.replace(doc, text=..., meta={**doc.meta, "lang": "en"}); it never assigns into doc.meta, because the caller may still hold the old document.

A generator is a function with yield in it. Calling it runs nothing; each next() runs it until the next yield. So

def length_filter(min_chars, max_chars):
def stage(docs):
for doc in docs:
if min_chars <= len(doc.text) <= max_chars:
yield doc
return stage

pulls one document from upstream, decides, yields it or not, and pulls the next only when downstream asks. compose(a, b, c) returns lambda docs: c(b(a(docs))): three generators stacked, still lazy. Taking 50 documents from the composed pipeline pulls exactly 50 from the source when no stage drops any. The only ways to break this are to collect (list(docs), sorted(docs), len(list(...))) or to read ahead in batches; either one makes memory grow with the corpus, and on an endless stream it never returns.

extract(manifest) is the first stage’s input: for each fetched or cached source of a data.01 manifest, in manifest order, it reads the raw parts with read_raw and yields Doc(id=f"{source_id}:{k}", ...) with meta = {url, fetched_at, license_spdx}; k restarts at 0 for each source. Quarantined sources are skipped by status.

The same visible text has many byte encodings. normalize_text maps them to one, in this order:

StepRuleWhy
1\r\n and lone \r become \nWindows and old Mac line endings
2remove every character of Unicode category Cc (control) except \n and \t, and remove U+FEFF (the byte order mark)\x00, \x07, escape sequences: invisible, never meant as text
3Unicode NFCé can be one code point (U+00E9) or e plus a combining acute accent (U+0301); NFC composes them, so equal text has equal bytes
4strip spaces and tabs at the end of each lineinvisible differences between copies
5runs of 3 or more \n become 2one blank line separates paragraphs (data.03 splits on it)
6strip white space at both ends

NFKC would also rewrite “compatibility” characters: the ligature fi becomes fi, ² becomes 2, full-width A becomes A. That changes what the text says, so this pipeline uses NFC and leaves the tokenizer to see what was written. Removing only category Cc keeps format characters (Cf) such as the zero-width joiner U+200D, which holds multi-person emoji together. The order matters for the law the tests check, idempotence (normalize_text(normalize_text(s)) == normalize_text(s)): if NFC ran before step 2, deleting a BOM between a and a combining accent would leave an uncomposed a plus accent that a second pass would compose. A document that normalizes to the empty string is dropped by normalize_unicode.

2.3 Language identification by stop-word share

Section titled “2.3 Language identification by stop-word share”

Stop words are the most frequent function words of a language (“the”, “of”, “and”); almost every sentence of ordinary prose holds several. filter.py carries a table STOPWORDS of 25 per language:

CodeStop words
deder die das und ist nicht ein eine zu den mit sich des auf für im dem von auch es er sie wir ich war
enthe of and to a in is it that was he she for on with as his her they at be this have from but
esel la los las de y que en un una es por con para del se no lo al su como pero muy yo está
frle la les de des et est un une du en que qui dans pour pas sur au ce il elle nous vous je avec
itil la di che e un una per non sono del della con si gli le ma come anche io questo nel è ha lo
nlde het een en van is dat op te in met niet zijn voor er maar ook als bij nog ik je wat hij ze
pto a os as de e que do da em um uma para com não no na por mais se ao dos das ele ela
SymbolMeaningType
w1,…,wNw_1, \dots, w_Nthe lang words of the text: maximal runs of letters in text.lower() (regex [^\W\d_]+)strings
SℓS_\ellSTOPWORDS[ℓ] for language code ℓ\ellset of strings
hℓ=∣{i:wi∈Sℓ}∣h_\ell = \lvert\{ i : w_i \in S_\ell \}\rverthits for ℓ\ell: how many words are stop words of ℓ\ellinteger
H=∑ℓhℓH = \sum_\ell h_\ellall hits; a word listed by two languages counts for bothinteger
ℓ∗=arg⁡max⁡ℓhℓ\ell^* = \arg\max_\ell h_\ellthe identified language; ties go to the alphabetically first codestring
c=hℓ∗/Hc = h_{\ell^*} / Hthe confidencefloat in (0,1](0, 1]

When H<2H < 2 there is too little evidence and the answer is ("und", 0.0) (undetermined). lang_filter(min_conf, langs) keeps a document when ℓ∗∈\ell^* \in langs and c≥c \ge min_conf, and records meta["lang"] and meta["lang_conf"] (data.06 writes them into the shard’s lang column). The corpus config’s default lang_min_conf = 0.65 is FineWeb’s threshold for fastText’s language classifier; here it is a share of votes, which is lower than fastText’s probability for the same text, so short documents that share many words with Dutch or Portuguese (“in”, “a”, “is”) can fall below it.

Rae et al. (2021, appendix A.1) filter a web corpus with cheap document statistics. With words = text.split() and nn = len(words):

Rule (name in GOPHER_RULES)A document fails whenDefault thresholdsCatches
wordsn<n < min_words or n>n > max_words50, 100,000fragments, dumps
mean_word_lenn=0n = 0, or the mean of len(w) is outside [min_mean_word_len, max_mean_word_len]3, 10character soup, run-together URLs
symbol_ratiotext.count("#") / n or (text.count("...") + text.count("…")) / n exceeds max_symbol_ratio0.1hashtag spam, teaser text
bullet_linesamong non-blank lines, the fraction whose lstrip() starts with •, ‣, ◦, -, or * exceeds max_bullet_lines0.9navigation menus, link lists
ellipsis_linesamong non-blank lines, the fraction whose rstrip() ends with ... or … exceeds max_ellipsis_lines0.3truncated previews
alpha_wordsn=0n = 0, or the fraction of words containing at least one letter is below min_alpha_words0.8tables of numbers, code dumps
stopwordsfewer than min_stopwords words, lowercased and stripped of the characters "'.,;:!?()[]{} at both ends, are in {the, be, to, of, and, that, have, with} (occurrences, not distinct words)2non-prose: keyword lists, boilerplate

Every boundary is inclusive: exactly 50 words, a ratio of exactly 0.1, exactly 30% ellipsis lines all pass. gopher_reasons(text, **thresholds) returns the names of the failed rules in the order above; gopher_rules(**thresholds) is the stage that keeps documents with no reasons, and it rejects an unknown threshold name when it is built, so a typo in the config fails before any document flows.

Spam and broken scrapes repeat themselves. Count how much of the text sits inside repeated word n-grams:

SymbolMeaningType
w0,…,wL−1w_0, \dots, w_{L-1}text.split()strings
gi=(wi,…,wi+n−1)g_i = (w_i, \dots, w_{i+n-1})the n-gram starting at word ii, 0≤i≤L−n0 \le i \le L - ntuple
duplicatedan n-gram value that occurs at least twice among the gig_i
coveredword wjw_j lies inside some occurrence gig_i of a duplicated n-gram (i≤j<i+ni \le j < i + n)
fn=∑j covered∣wj∣∑j∣wj∣f_n = \dfrac{\sum_{j \text{ covered}} \lvert w_j \rvert}{\sum_j \lvert w_j \rvert}dup_ngram_char_frac(text, n): the share of characters (spaces not counted) inside repeated n-grams; 0 when L<nL < nfloat in [0,1][0, 1]

Every occurrence of a duplicated n-gram counts, the first one included. repetition_filter(n, max_frac) keeps documents with fn≤f_n \le max_frac (config defaults n=3n = 3, 0.3).

2.6 Perplexity and the configured pipeline

Section titled “2.6 Perplexity and the configured pipeline”

ppl_filter(lm, max_ppl) keeps documents a language model does not find too surprising: any object with perplexity(text) -> float works, and in C1 the L2.1 n-gram model trained on clean text plugs in. The score is recorded as meta["ppl"].

build_filters(filters, lm=None) turns the corpus config’s filters table into one stage, defaults from the schema, in this order: normalize_unicode, length_filter(min_chars, max_chars), lang_filter(lang_min_conf, langs), gopher_rules(**gopher), repetition_filter(repetition.n, repetition.max_frac), and ppl_filter(lm, ppl.max_ppl) when the table has ppl.max_ppl (a ValueError when it does but no lm is passed). Normalization comes first so every later rule counts the characters the model will see.

Normalize. Take the 41 code points

U+FEFF Z o e U+0308 " and the cat" \r\n \r\n \r\n \r\n "sat on a mat." \t " " U+0007

Step 1 turns the four \r\n into four \n. Step 2 removes U+FEFF and the BEL control U+0007. Step 3 composes e + U+0308 into ë (U+00EB). Step 4 strips the tab and space at the end of the last line. Step 5 turns the four \n into two. The result is Zoë and the cat\n\nsat on a mat., and normalizing it again changes nothing.

Identify. In “The cat and the dog sat in a box.” the lang words are the, cat, and, the, dog, sat, in, a, box:

Worddeenesfritnlpt
the (twice)2
and1
in11
a11
total hℓh_\ell0500011

H=7H = 7, ℓ∗=\ell^* = en, c=5/7≈0.714≥0.65c = 5/7 \approx 0.714 \ge 0.65: kept, with meta["lang"] = "en".

Gopher. “Buy now. #deal #sale #wow …” has 6 words: words fails (6 < 50); the mean word length is (3+4+5+5+4+3)/6=4(3 + 4 + 5 + 5 + 4 + 3)/6 = 4, fine; symbol_ratio fails (3/6=0.5>0.13/6 = 0.5 > 0.1, and the ellipsis alone gives 1/61/6); its one line ends with ..., so ellipsis_lines fails (1/1>0.31/1 > 0.3); 5 of 6 words hold a letter (0.83≥0.80.83 \ge 0.8); no Gopher stop word appears, so stopwords fails. gopher_reasons returns ["words", "symbol_ratio", "ellipsis_lines", "stopwords"].

Repetition. “the cat sat the cat sat down”, n=3n = 3: the trigrams are (the cat sat), (cat sat the), (sat the cat), (the cat sat), (cat sat down). Only (the cat sat) occurs twice, at i=0i = 0 and i=3i = 3, so words 0 to 5 are covered: 6×3=186 \times 3 = 18 characters of 18+4=2218 + 4 = 22, f3=18/22≈0.818f_3 = 18/22 \approx 0.818.

@dataclass(frozen=True, slots=True)
class Doc: id: str; source_id: str; text: str; meta: Mapping[str, Any] = {}
Stage = Callable[[Iterator[Doc]], Iterator[Doc]]
def compose(*stages: Stage) -> Stage: ...
def extract(manifest: Manifest) -> Iterator[Doc]: ...
def normalize_text(text: str) -> str: ...
def normalize_unicode(docs: Iterator[Doc]) -> Iterator[Doc]: ... # a Stage itself
def length_filter(min_chars: int = 50, max_chars: int = 100_000) -> Stage: ...
def identify_lang(text: str) -> tuple[str, float]: ...
def lang_filter(min_conf: float = 0.65, langs: Sequence[str] = ("en",)) -> Stage: ...
def gopher_reasons(text: str, *, min_words=50, max_words=100_000, min_mean_word_len=3.0,
max_mean_word_len=10.0, max_symbol_ratio=0.1, max_bullet_lines=0.9,
max_ellipsis_lines=0.3, min_alpha_words=0.8, min_stopwords=2) -> list[str]: ...
def gopher_rules(**thresholds: float) -> Stage: ...
def dup_ngram_char_frac(text: str, n: int) -> float: ...
def repetition_filter(n: int = 3, max_frac: float = 0.3) -> Stage: ...
def ppl_filter(lm: Any, max_ppl: float) -> Stage: ...
def build_filters(filters: Mapping[str, Any], lm: Any = None) -> Stage: ...

STOPWORDS, GOPHER_STOPWORDS, GOPHER_RULES, BULLETS, and ELLIPSES are module constants with the values of section 2. Standard library only (re, unicodedata, collections, dataclasses); extract reads through corpus.fetch.read_raw.

TestKINDChecksWhy it matters downstream
test_hand_example_normalize_and_identifyunitthe section 3 normalization and ("en", 5/7)you and the test agree on both definitions
test_normalize_rules_one_by_oneunitone hand-labelled case per rule of 2.2, including NFC not NFKC and a kept ZWJeach rule is what the corpus will look like
test_normalize_is_idempotentproperty500 seeded random strings over the rules’ special charactersreruns and re-normalized text never drift
test_normalize_unicode_drops_empty_and_keeps_identityunitan empty-after-normalization document is dropped; id and meta keptno empty rows in shards
test_length_filter_bounds_are_inclusiveboundarymin_chars and max_chars both passthe config’s numbers mean what they say
test_identify_lang_hand_labelledunitEnglish, German, French, Dutch, digits, one stop wordthe lang column of the shards
test_identify_lang_ties_go_to_the_first_codeboundary“de de” is ("es", 0.25)the answer never depends on dict order
test_identify_lang_uses_the_stopword_tablepropertyconfidence is k/(k+j)k/(k + j) for your own tablethe formula, separate from the word lists
test_lang_filter_threshold_is_inclusive_and_records_metaboundaryc=c = min_conf is kept; meta set on a new document, input untouchedstages share documents safely
test_gopher_rules_hand_labelledunitone text failing exactly each rule, a clean story, punctuated words, the empty texteach rule catches its own junk
test_gopher_boundaries_are_the_papersboundary50 words, ratio 0.1, 30% ellipsis lines, 90% bullets all passreal text at every boundary survives
test_gopher_stopwords_count_occurrences_after_stripping_punctuationboundary“The,” and “(the)” count; one occurrence is not enoughprose with punctuation passes
test_gopher_rules_reject_unknown_thresholds_when_builtboundarymin_word=10 raises when the stage is builtconfig typos fail early
test_dup_ngram_hand_exampleunitthe section 3 value 18/2218/22 and edge casesthe repetition statistic itself
test_repetition_filter_keeps_the_boundaryboundaryfn=f_n = max_frac is kept; n = 0 rejectedthe threshold is inclusive
test_ppl_filter_takes_any_perplexity_objectunita fake model; equal to max_ppl kept; meta["ppl"] setC1 plugs in the L2.1 model
test_compose_runs_left_to_rightunitcompose(a, b) is b after a; compose() is the identitynormalization runs before counting
test_stages_are_lazyproperty50 documents taken from an endless source pull exactly 50no stage collects its input
test_pipeline_memory_stays_flatproperty2,000 documents through the whole pipeline peak under 0.5 MB of allocationsthe pipeline fits an activity’s memory limit
test_golden_200_documentsconformance200 generated documents: kept ids, order, lang, and the sha256 of every output textthe configured pipeline, byte for byte
test_build_filters_validates_the_tableboundaryunknown key, unknown Gopher name, ppl without a model; ppl with a modelreviewed config, early errors
test_gopher_rule_names_are_the_contractunitGOPHER_RULES names and orderthe config and the datasheet use these names
test_extract_numbers_documents_per_source_across_partsunitids a:0 to a:2 across two parts, then b:0; meta from the raw linestable ids for every later stage
test_extract_skips_quarantined_sourcesboundarya quarantined entry contributes nothing, even with parts on diskbytes that failed their checksum never enter

The golden corpus (course/oracle/data.02/golden_corpus.py) has 125 documents to keep, 15 of them noisy copies of good stories (CRLF, NFD accents, a BOM, controls, trailing blanks) whose expected output is the clean story, and 75 to drop, each built to fail one named rule by a wide margin: 5 empty, 5 too short in characters, 15 in German, French, Spanish, or Dutch, 12 under 50 words, 8 hashtag-heavy, 8 bullet lists, 7 trailing off in ellipses, 7 a third numbers, and 8 repeating one sentence six times.

PitfallSymptomCaught by
1. NFKC instead of NFCligatures, superscripts, and full-width text silently rewrittentest_normalize_rules_one_by_one (mutant s01)
2. Removing every invisible character (category Cf too)emoji sequences joined by U+200D fall aparttest_normalize_rules_one_by_one (mutant s03)
3. Composing (NFC) before removing controlsnormalizing twice gives a different texttest_normalize_is_idempotent (mutant s05)
4. conf > min_confdocuments exactly at the threshold droppedtest_lang_filter_threshold_is_inclusive_and_records_meta (mutant s09)
5. doc.meta["lang"] = ... on the incoming documentan upstream holder of the document sees meta it never settest_lang_filter_threshold_is_inclusive_and_records_meta (mutant s10)
6. w.isalpha() for “contains a letter”every punctuated word (“cat,”) counts as non-alphabetic; ordinary prose failstest_gopher_rules_hand_labelled (mutant s14)
7. Comparing stop words without lowercasing and stripping punctuation“The,” does not count; short clean documents fail stopwordstest_gopher_stopwords_count_occurrences_after_stripping_punctuation (mutant s19)
8. Ignoring unknown threshold namesa typo in the config silently uses the defaulttest_gopher_rules_reject_unknown_thresholds_when_built (mutant s20)
9. A stage that reads its input in batches, or calls list()memory grows with the corpus; on an endless stream it never returnstest_stages_are_lazy (mutant s26)
10. A stage that keeps every document it has seenthe activity runs out of memory halfway through the corpustest_pipeline_memory_stays_flat (mutant s32)
11. Counting only the repeats after the first occurrencerepetition is underestimated by half; spam passestest_dup_ngram_hand_example (mutant s21)
DirectionModuleHow it uses this
Backdata.01extract reads the raw parts through read_raw and skips quarantined sources from the fetch Manifest
Forwarddata.03exact_dedup is a stage over your Docs, run after the filters
Forwarddata.04near dedup (MinHash, LSH) is a stage over the same Docs, composed after exact dedup
Forwarddata.05the PII scrub is a stage that returns new Docs with redaction counts in meta
Forwarddata.06the shard writer consumes the final Doc stream and writes meta["lang"] into the lang column
ForwardC1ppl_filter with the L2.1 n-gram model trained on clean text

If you skip this module, ol check data.03 stops with data.03 needs data.02: build it, or rerun with --ref-deps.

Your pieceProduction equivalentWhat it addsWhere to look
Stage, composedatatrove PipelineStep and executorsthe same iterator-in, iterator-out steps, sharded across processes and Slurm tasks, with per-step statsdatatrove/pipeline/base.py
gopher_rulesdatatrove GopherQualityFilter; FineWeb’s custom filtersthe full Gopher set plus FineWeb’s line-level rules (lines ending in punctuation, short-line ratio)datatrove/pipeline/filters/gopher_quality_filter.py
repetition_filterdatatrove GopherRepetitionFilterduplicate lines, paragraphs, top n-gram fractions for n = 2 to 4, duplicate n-gram fractions for n = 5 to 10datatrove/pipeline/filters/gopher_repetition_filter.py
identify_langfastText lid.176, GlotLID176 (or 1,600+) languages from character n-grams, calibrated probabilitiesfasttext.cc, “Language identification”
ppl_filterCCNet’s KenLM perplexity bucketsa 5-gram model per language, documents bucketed into head, middle, tailWenzek et al. 2019, CCNet