TT Lab
Get started
Learn Learning paths Courses

Data Pipelines

Lineage and Reproducibility: Run It Twice, Compare the Bytes

Continue in TT Lab

Goal

You build a transformer lineage.py that aggregates the order drops that fall every day. You leave in a manifest which inputs (content hashes), which code, and which parameters each output was made with, prove by running twice with the same input that the bytes are the same, and confirm by experiment which input columns a column came from.

Why it matters

When a report number is off from accounting, the question you must answer is one — where did this number come from. You cannot answer with file names and modification times. It is common for the name to be the same while upstream overwrote the file, and the modification time changes from just copying. The only thing that can serve as evidence is the hash of the content itself. And for "let's rerun it" to work, the same input must give the same output. That does not happen by itself. The current time written into the output, an identifier newly drawn on every run, an unsorted file list and set iteration, and the order in which floating-point numbers are added quietly break reproduction. This lab finds and removes those four one by one. Lineage also has a granularity. The dataset level goes as far as "this table came from those three files," and the column level goes as far as "this column came from those input columns." The granularity at which you can answer when upstream tells you they are changing one column is the latter. And column lineage goes stale as soon as you write it down, so you remeasure it by experiment. The grader does not trust the text you write down. It sets up the drops that the grader created in a temporary directory, actually runs your transformer, compares the aggregate values and the manifest against values the grader counted directly, runs twice to compare the bytes, and shakes the input columns one at a time to check that the column lineage you wrote matches. The shop names and amounts change on every run.

Steps

  1. Create and run /root/lineage/gen_drops.py to create three date-based drops under /root/lineage/drops.
  2. In /root/lineage/lineage.py, create digest so that it produces the content hash, size, and record count.
  3. Add run to produce the aggregate output /root/lineage/out/shops.csv.
  4. With run --manifest, also produce the manifest /root/lineage/out/shops.manifest.json.
  5. Fix it so that running twice gives the same bytes, and leave the result in /root/lineage/repro.json.
  6. Add columns to produce column-level lineage and leave it in /root/lineage/columns.json.
  7. Add verify to rebuild according to the manifest and compare.
  8. Add trace and write /root/lineage/lineage_report.md.

Reference

Create the drops that fall every day

Create and run /root/lineage/gen_drops.py to create three date-based drops under /root/lineage/drops. The header is order_id,shop,qty,amount,dropped_at, and there must be a mix of lines with quantity 0 and lines where the same voucher came again with only the amount changed.

In the field, the files come first. Here we create those files. There must be at least four shops, at least three accepted lines per shop, and one voucher number must appear in two different files. dropped_at is a column the aggregation does not use, so it is enough for it to have several different values.

Judge sameness by content, not by name

In /root/lineage/lineage.py, create digest <파일> (the placeholder stands for the file) so that it produces path, sha256, bytes, and rows as JSON. rows is the record count excluding the header.

Read the file in binary chunks and feed them to hashlib.sha256. Copy the same content under a different name, or change only the modification time with touch, and confirm for yourself that the hash stays the same. You must count records with csv, not by the number of lines.

Produce the aggregate output

Add run --drops <디렉터리> --out <파일> [--min-qty N] (the placeholders stand for the directory and the file) to produce /root/lineage/out/shops.csv. Read the files in name order, keep only the first-appearing line for the same voucher, then discard lines whose qty is less than min_qty, group by shop, and write shop,orders,qty,amount_cents in ascending shop order.

The fact that deduplication comes before the quantity filter changes the answer — if the first-appearing line has quantity 0, that voucher drops out entirely. Convert the amount from a string to integer cents and add. Read the file list sorted. The order of os.listdir is not promised.

Produce the output and the manifest together

Add run --manifest <파일> (the placeholder stands for the file) to also produce /root/lineage/out/shops.manifest.json. It holds the eight keys tool, code_sha256, drops_dir, params, inputs, output, run_id, and created_at.

For inputs, put in a name-ordered list holding name, sha256, bytes, and rows, and code_sha256 is the hash of lineage.py itself. Write the defaults for the parameters too — defaults change later. If you attach it later, it is certainly missed, so create it inside the same function as the output.

Is it the same bytes when run twice

Fix it so that, when run twice with the same input, the bytes of the output are the same and the manifest is the same in every field except created_at. And leave the results of running twice in /root/lineage/repro.json as run_a, run_b, identical, and manifest_diff_keys.

The current version draws the run identifier from random numbers, so the manifest is different every time. Derive the identifier from the content — if you concatenate the code hash, the parameters, and the output hash and hash them, the same run has the same name. Also check again whether you sorted the file list and the group list. Python has a different string-hash seed on each run, so the set iteration order changes between runs.

Measure by experiment where a column came from

Add columns to produce, for each output column, the sorted list of input columns it depends on, and leave the same content in /root/lineage/columns.json.

Do not stop at writing it down; confirm it by experiment — shake one input column, rerun, and see which output column changes. If you change the shop, the output's lines themselves change, so every column depends on it. If you make a voucher number the same as another line's, that line drops out as a duplicate, so the count, quantity, and amount move together. A column the aggregation does not read at all must appear nowhere.

Rebuild according to the manifest and compare

Add verify <매니페스트> (the placeholder stands for the manifest) to compare the input hashes and the code hash with the current values, and to rebuild using the manifest's inputs and parameters and compare the output hash. If even one does not match, the exit code is 4.

For inputs, report separately those that are gone and those whose content changed — they are different stories for the person investigating. You must rerun with the parameters written in the manifest. If you run with today's defaults, even when a different answer comes out from the one then, you cannot pinpoint the cause.

Answer which input that number came from

Add trace --manifest <파일> --key <상호> (the placeholders stand for the file and the shop) to produce that shop's count, quantity, and amount, and the input files that actually added its lines, by name, content hash, and number of lines. And write /root/lineage/lineage_report.md with four sections: ## 이 숫자는 어디서 왔나 ## 이름과 시각은 왜 근거가 못 되나 ## 두 번 돌려 같았는가 ## 칼럼 하나가 어디서 왔나 (in order: where this number came from, why names and times cannot serve as evidence, whether it was the same when run twice, and where a column came from).

Put in sources only the files that actually added lines for that shop — a file that was read but added not a single line is not lineage. Use the hash as written in the manifest. In the report, write numbers; the person investigating must be able to open that line in their own file.