Data Engineering
Build a complete local evidence pipeline first, then extend its contracts, replay, quality, history, authority, cost, and recovery proofs to analytical and distributed adapters.
Data contracts and reproducible ingestion
Objective Run one complete local ingestion and prove the grain, identity, time rule, per-record outcome, conservation total, and reproducible manifest.
Core explanation
A data pipeline moves facts from a source into a product while preserving their meaning and evidence. A data contract states what one record represents, its stable identity, required fields, value rules, time meaning, producer, owner, version, and allowed change. In the first laboratory, one line represents one recorded learner-score event and event_id is its producer-stable identity. The source is newline-delimited JSON, meaning each line is one complete JSON object. Ingestion reads every line without changing the source file, classifies it, and writes accepted, duplicate, or quarantined evidence to separate outputs. Quarantine means retain an unusable record with its source location and stable reason so it can be inspected or repaired; it does not mean silently discard it. A fixed run-as-of time makes the same source and code produce the same time decision on every replay. Conservation requires source count to equal accepted plus duplicate plus quarantined counts. A manifest records contract, code label, source hash, run parameters, counts, and output hashes so another operator can identify exactly what was processed. This standard-library path proves local file parsing, classification, conservation, and reproducibility under the recorded conditions; it does not prove a connector, object store, distributed engine, scheduler, stream, warehouse, or production recovery. Later chapters add replay, modeling, layout, orchestration, quality, CDC, streaming, table commits, governance, operations, and recovery only after their terms and evidence are introduced.
Run one complete ingestion from a source file to an evidence directory
Create an empty `ingestion-boundary` directory. Split the example at its two file markers: save the four JSON lines under `# source-events.ndjson` as `source-events.ndjson`, and save everything after `# ingest.py` as `ingest.py`. Run `python3 ingest.py` on macOS or Linux or `py ingest.py` on Windows. The program uses only the Python standard library and works inside the current directory; no package installation, service, database server, scheduler, or account is involved.
Expect `source=4 accepted=1 duplicates=1 quarantined=2` and `conserved=True`. Open `run-001` and confirm four files exist: accepted, duplicate, and quarantine NDJSON plus `manifest.json`. Do not begin by reading every function. First connect each visible artifact to one question: which source facts were accepted, which repeated, which could not be trusted, and which exact contract, input bytes, parameters, code label, counts, and output bytes produced this run?
Define record grain, stable identity, and three clocks
Grain completes the sentence “one record represents one …”. Here one line represents one recorded learner-score event. It is not one learner, one current score, or one file delivery. `event_id` is the producer-stable event identity: an identical redelivery of e1 is the same logical event, while a different payload under e1 is a conflict that needs human or producer resolution. A source line locates evidence but is not business identity because ordering or file boundaries can change during transport.
`occurred_at` is event time: when the learning event happened. Source arrival time would describe when the pipeline received the file, and processing time would describe when this run evaluated it. Those clocks answer different questions. The program uses the fixed `RUN_AS_OF` processing boundary so a replay next month does not silently reclassify a future event. The contract also names required fields, score range, time-zone requirement, and version. A repair should create a new source or contract version rather than silently rewriting preserved evidence.
Give every source line one preserved and counted outcome
Line 1 is accepted because it has the required fields, valid aware timestamp, allowed score, and unseen identity. Line 2 is quarantined with `invalid_time`; the original parsed record and source line remain inspectable. Line 3 is an identical e1 redelivery, so it goes to the duplicate evidence file. Line 4 is quarantined with `score_out_of_range`. Quarantine is not deletion and duplicate is not automatically corruption. Each category has a stable reason and source location so later repair or audit can begin from evidence.
Conservation is the arithmetic boundary `source = accepted + duplicates + quarantined`, here `4 = 1 + 1 + 2`. The manifest binds that result to a contract, fixed code label, run-as-of parameter, SHA-256 source identity, and output identities. Sorting object keys and avoiding current time make the written bytes deterministic for this fixture. Run the unchanged program twice and the source, outputs, counts, and manifest should remain byte-for-byte identical. That is local reproducibility; it does not yet prove retry-safe publication or distributed execution.
Use one value mutation and one identity conflict to test the contract
Change only e3’s score from 120 to 90. Predict `source=4 accepted=2 duplicates=1 quarantined=1` with conservation still true. Run twice, inspect the moved e3 evidence and stable manifest, then restore 120 and the 4/1/1/2 baseline. This mutation tests the score rule without changing grain, identity, or time. If another count changes, locate the first classification branch whose promise became false.
Next change only the repeated e1 score from 82 to 83. Predict one accepted, zero duplicates, and three quarantined because the repeated identity now has a conflicting payload. Verify `conflicting_duplicate` is preserved in quarantine, restore 82, and recover the baseline. Give the two files and evidence directory to another learner. They must reproduce both experiments and explain why dropping the conflicting line, accepting “latest wins,” or changing the original source in place would destroy evidence rather than resolve the contract.