Skip to content

Parquet shards, manifest, and the document-hash split

Moduledata.06 · build · Python · Pass 3 · 3 to 4 h
You buildpython/corpus/shard.py: split_of, shard_name, write_shards, read_shards, and the tables SCHEMA, DEDUP_DEFAULTS, DEFAULT_ROW_GROUP_BYTES
Contractcourse/contracts/py/corpus/shard.pyi · files: formats/corpus-shard.md, formats/corpus-manifest.schema.json · allowed library: pyarrow (allowed-deps.toml)
Testscourse/tests/data.06/ (what they check: section 4)
Needsdata.02 Doc · data.04 the minhash_cluster tag · data.05 the PII counts and manifest_counts (or --ref-deps) · reading: Storage & Warehousing (columnar files)
Used bydata.07 tokenizes the shards through read_shards · data.08 joins every row with the ledger
MilestoneMS-corpus
Optional depthApache Parquet file format specification (free); Abadi et al., The Design and Implementation of Modern Column-Oriented Database Systems (2013), chapters 1 to 3 (free); the FineWeb datatrove writers (going further)
  • The corpus becomes files with a contract: Parquet shards in a fixed column order and type, sorted by id, cut every shard_rows rows, and a manifest that records each file’s hash and every count the pipeline dropped (test_schema_compression_and_manifest).
  • A document’s split is a function of its text alone, the first 8 bytes of its SHA-256 modulo 1000, so exact copies can never straddle train and val and the split never moves when documents are added or reordered (test_split_follows_the_hash).
  • The same documents in any order give byte-identical shards and manifest (test_output_is_deterministic), and a failed write never damages the previous output (test_atomic_replace).
Terminal window
ol start data.06 # stubs shard.py into python/corpus/
ol tests data.06 # the course test catalog
ol tdd red data.06 # rung R3: your tests first, failing against the stubs
ol check data.06 # exit code is the verdict
ol check data.06 --ref-deps # only if data.02, data.04, or data.05 is not passing yet
ol mutate data.06 # how many planted bugs your tests catch
ol diff data.06 # after passing: your code against the reference

Add pyarrow to your python/pyproject.toml (dependencies = ["numpy>=2", "pyarrow>=18", "zstandard>=0.23"]); the contract pre-check allows it for corpus units only.


After data.05 the cleaned corpus exists only as a Python iterator. Nothing downstream can use it: the tokenizer (data.07) needs files it can read in a fixed order, the training loader needs a validation split that never overlaps training, the ledger check (data.08) needs every row’s source and license, and the durable CorpusBuild (data.09) needs an output whose hash says whether a retried run produced the same thing. A split made by random() would move documents between train and val on every run, and an exact copy that sits in both makes your validation loss lie. This module writes the corpus as sorted, hash-split Parquet shards with a manifest, atomically.

SymbolMeaningType / shape
NNnumber of documentsint
RRrows per shard (shard_rows)int
n=max⁡(1,⌈N/R⌉)n = \max(1, \lceil N / R \rceil)number of shardsint
h(d)h(d)SHA-256 of the UTF-8 bytes of document dd‘s text32 bytes
u(d)u(d)the first 8 bytes of h(d)h(d) as a big-endian unsigned integerint in [0,264)[0, 2^{64})
vvval_permille: the share of val documents in thousandthsint in [0,1000][0, 1000]
BBrow_group_bytes, the text budget of a row groupint

Columnar files. A Parquet file stores each column contiguously, compressed (zstd here), in row groups: blocks of rows that a reader loads together. Reading only source_id and license_spdx (what data.08 needs) touches a fraction of the bytes; reading text in order (what data.07 needs) streams one row group at a time, so BB bounds a reader’s memory. The schema is fixed in formats/corpus-shard.md: id, text, source_id, url, license_spdx, lang, n_chars (int32), sha256 (fixed_size_binary[32]), minhash_cluster (int64), pii_redactions (int32), split, in that order.

Sort, then cut. All rows are sorted by id as strings (Python’s str order, which equals UTF-8 byte order): "a:10" comes before "a:2". Shard ii holds rows iRiR to (i+1)R−1(i+1)R - 1 and is named shard-{i:05d}-of-{n:05d}.parquet; an empty corpus still has one empty shard, so a reader never finds a manifest without a file. Inside a shard, consecutive rows go into one row group while their text bytes sum to at most BB; a longer document gets a group of its own.

The document-hash split. A document is val when

u(d) mod 1000<v,u(d) \bmod 1000 < v,

else train. The hash depends only on the text, so the decision is the same in every run, on every machine, in any order, and for every copy of the same text. Over many documents u mod 1000u \bmod 1000 is close to uniform, so the val share is close to v/1000v / 1000: with 4000 documents and v=100v = 100, the val count is Binomial(4000, 0.1), mean 400, standard deviation 19.

Columns from upstream. minhash_cluster: rows sharing meta["minhash_cluster"] (the root id data.04 wrote; a document without it is its own cluster) form a cluster, and its value is the smallest row index of the cluster in output order, a number you can group by without the ids. pii_redactions is meta["pii_redactions"]; the manifest’s pii object sums data.05’s manifest_counts(meta["pii"]) over the documents. lang defaults to und (undetermined), n_chars counts code points. A document without url or license_spdx is an error: the ledger could not trace it.

The manifest. _MANIFEST.json lists every shard with its file’s SHA-256 and row count, the dropped counts of every stage (filters as given, dedup as DEDUP_DEFAULTS overridden by the caller), the pii totals, the ledger path, and the config’s hash, keys in the format’s order, written as json.dumps(m, indent=2, ensure_ascii=False) plus a newline. Two writers that follow the contract write the same bytes.

Atomic output. Write everything into out.tmp (removing a stale one first), then rename it to out, replacing an older out. A crash before the rename leaves the previous output untouched; a crash after it leaves the new one complete. A rerun of the durable activity (data.09) sees either the old or the new version, never half of each.

Three documents, val_permille 300, shard_rows 2. web:2 is a near duplicate of web:10, kept with drop=False, so both carry minhash_cluster = "web:10":

idtextfirst 8 bytes of SHA-256as an integer, mod 1000split
web:2Tom has a red ball.ba7f5216fe451c8113438550071805549697 → 697train
web:10Tom has a red ball!39777282a7bff2424140904287876149826 → 826train
books:0Mia naps.fcd7f3a1b48b415818219298693394940248 → 248val (248 < 300)

Sort. By id as strings: books:0 < web:10 < web:2 (b before w; then 1 before 2, character by character). Rows 0, 1, 2.

Clusters. web:10 and web:2 share the tag, and their smallest row is 1; books:0 is alone at row 0. Column: [0, 1, 1]. n_chars: [9, 19, 19].

Shards. n=⌈3/2⌉=2n = \lceil 3 / 2 \rceil = 2: shard-00000-of-00002.parquet holds rows 0 and 1, shard-00001-of-00002.parquet holds row 2.

The boundary. split_of(h("Mia naps."), 248) is train: 248<248248 < 248 is false. That strict inequality is what makes v=0v = 0 an empty val split and v=1000v = 1000 all val.

This is test_hand_example; read_shards gives the same three rows back in test_read_shards.

# python/corpus/shard.py (the full contract is contracts/py/corpus/shard.pyi)
SCHEMA: pa.Schema
DEFAULT_ROW_GROUP_BYTES: int # 64 MiB
DEDUP_DEFAULTS: dict[str, Any]
def split_of(sha256: bytes, val_permille: int) -> str
def shard_name(i: int, n: int) -> str
def write_shards(docs, out: Path, *, dataset, version, shard_rows, val_permille=5,
row_group_bytes=DEFAULT_ROW_GROUP_BYTES, filters=None, dedup=None,
ledger_ref="corpus/LEDGER.jsonl", config_sha256="0" * 64) -> dict
def read_shards(out: Path, split: str | None = None) -> Iterator[Doc]

write_shards has to read every document before writing (it sorts), so it is the end of the streaming part of the pipeline. Write with pq.ParquetWriter(path, SCHEMA, compression="zstd") and one write_table per row group (row_group_size set to the group’s length, so pyarrow does not split it again).

TestKINDChecksWhy it matters downstream
test_hand_exampleunitthe section 3 order, splits, clusters, character counts, and shardsyou and the tests agree on the definitions
test_schema_compression_and_manifestconformanceexact Arrow schema, zstd, manifest valid against the schema, file hashes and row counts true, no stray filesevery reader relies on the contract
test_sorted_by_id_in_byte_orderunitstring order, not numeric; n_chars in code points; und default; untagged rows are their own clusterrows land in the same shard on every machine
test_split_follows_the_hashpropertythe big-endian rule on every row; exact copies on one side; split stable under reordering and subsetsno val text leaks into training
test_val_share_matches_permillestatistical4000 documents at 100 permille: val within 4 standard deviations of 400; bad arguments raisethe configured share is what you get
test_shards_and_namesboundary23 rows at 10 per shard: 10, 10, 3; names; an empty corpus has one empty shardreaders can always open shard 0
test_row_groups_are_cut_by_text_bytesboundarygroups fit the byte budget; a long document gets its own group, also first in a shardbounded reader memory
test_output_is_deterministicpropertyshuffled input gives byte-identical filesMS-corpus and data.09 compare hashes
test_cluster_column_from_near_dedupunitdata.04’s tags become the smallest row indexthe column analysts group by
test_pii_totalsunitdata.05’s counts summed into the manifest; the per-row columnthe datasheet’s numbers
test_counts_and_key_orderunitfilters as given, dedup over the defaults, the format’s key order and serializationtwo writers, one manifest
test_rejects_bad_inputboundarymissing url or license, repeated id, shard_rows 0: errors, nothing writtenthe ledger can trace every row
test_atomic_replacefaulta failed run leaves the old output; a stale .tmp is clearedretried activities are safe
test_read_shardsunitmanifest and row order, split filter, meta columns, a rewritten shard raisesdata.07 and data.08 read through it

Your tests (rung R3, red then green). Under python/tests/data-06-shard/, failing first against the stubs (ol tdd red data.06): the hand example; the split as a function of the text (with the big-endian rule written out); the 248 boundary; shard counts and names, empty corpus included; string order of ids; a missing license raising; a stale .tmp cleared on success; PII totals and file hashes; same input, same bytes; column defaults; recorded counts; repeated ids and a tampered shard. ol mutate data.06 grades them: 0.70 of the mutants, including the one behind Pitfall 2.

PitfallSymptomCaught by
1. keeping arrival order, or sorting ids numericallyshards differ between runs, or from every other implementationtest_sorted_by_id_in_byte_order, test_output_is_deterministic (mutants s01, s08)
2. splitting by the id, or by a random drawan exact copy under another id lands in val while the original trainstest_split_follows_the_hash (mutant s02)
3. cutting a row group only after it overflowsgroups exceed the budget; a reader’s memory is unboundedtest_row_groups_are_cut_by_text_bytes (mutant s10)
4. defaulting a missing license to ""rows the ledger cannot trace reach trainingtest_rejects_bad_input (mutant s12)
5. writing straight into outa crash leaves half a corpus that looks completetest_atomic_replace (mutant s13)
6. reading the first 8 bytes little-endiana different, still uniform split: every other implementation disagreestest_split_follows_the_hash (mutant s03)
7. a timestamp in the manifestthe output hash changes on every runtest_output_is_deterministic (mutant s11)
8. trusting the shard bytes in read_shardsa rewritten shard feeds the tokenizer silentlytest_read_shards (mutant s14)
DirectionModuleHow it uses this
Backdata.02Doc, and meta["lang"] from the language filter
Backdata.04meta["minhash_cluster"] becomes the row-index column
Backdata.05pii_redactions, and the per-kind counts renamed by manifest_counts
Forwarddata.07read_shards(out, split) feeds the tokenizer, train then val, in manifest order
Forwarddata.08every row’s source_id and license_spdx are checked against the ledger; the datasheet reads the manifest
Forwarddata.09CorpusBuild’s shard activity; its output hash is the retry check

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

Your pieceProduction equivalentWhat it addsWhere to look
write_shardsdatatrove ParquetWriter, Hugging Face datasetsstreaming writers that never hold the corpus, file rotation by size, remote filesystemsdatatrove src/datatrove/pipeline/writers/parquet.py
the manifestDelta Lake and Apache Icebergtransaction logs with snapshots, schema evolution, time travelthe Iceberg table spec
split_ofHugging Face datasets train/test splits, TFDSdeterministic splits by hash of a key; named splits in metadatatensorflow_datasets/core/splits.py
atomic renameobject-store commit protocolsmulti-file commits without rename (S3 has none): manifest files written lastthe Iceberg commit protocol