TT Lab
Get started
Learn Learning paths Courses

Data Pipelines

The 3 A.M. Settlement Job That Died: A Ledger and an Atomic Swap

Continue in TT Lab

Goal

You build a runner runner.py that is safe even if it dies midway. It writes outputs under a temporary name and swaps them in so that no half-written file remains, leaves one run as one line in a ledger, resumes from the shard that died, and keeps the numbers from growing even if the same run is committed twice.

Why it matters

A pipeline certainly dies midway. The problem is not the fact that it dies but what remains when it does. If it was writing directly to the destination file, a half-written file remains, and that file has a perfectly fine size and name, so the next step reads it as is. Running it again is also dangerous. Without a record of how far the previous run got, there is no choice but to run from the start, and at the last place where you append to the ledger, the same amount is added twice. That is the real identity of "why did it double when I reran it." This lab builds three devices. The first is atomic replacement. You write everything under a temporary name in the same directory and then swap it in with os.replace. The second is the run ledger. You append one line per run, and if it failed, you write which shard it died at. The third is resume. You use the output of a finished shard itself as the mark and skip it. The dp-idempotent lab in this course deals with idempotency on the database side — making the same row become one row even if inserted twice. This one comes before that. It deals with what remains in the filesystem at the place where the process died, and, looking at what remains, from where to start again. The grader does not trust the text you write down. It sets up input shards that the grader created in a temporary directory and actually runs your runner. After deliberately killing it, it even checks that the destination file is intact, that the temporary file remains next to the destination, and that the name of the failed shard is written in the ledger. The number of shards and the amounts change on every run.

Steps

  1. Create and run /root/runx/gen_shards.py to create the input shards under /root/runx/work/in.
  2. In /root/runx/runner.py, create scan so that it produces the shard list, the count, and the total.
  3. Add part to process one shard, write the output under a temporary name, and swap it in. Also build the path that dies just before the swap with --crash=write.
  4. Add run to process all the shards and produce the merged output.
  5. Leave one ledger line per run and produce a summary with ledger.
  6. Make it die at a middle shard with --crash-shard, and check that the failed shard remains in the ledger and the merged output is not touched.
  7. Add --resume to skip finished shards and continue from where it died.
  8. Add commit so that the same run is not attached twice to the day's ledger /root/runx/work/out/daily.jsonl.

Reference

Create the shards the upstream dropped

Create and run /root/runx/gen_shards.py to create shard files under /root/runx/work/in. There must be at least 4 shards, at least 5 lines per shard, and at least 40 lines in total, and one line is a JSON object holding an id and an integer amount.

One shard is one JSON Lines file. The file name with .jsonl removed becomes the shard name, so pad the positions like h00 and h01 so that sorting gives time order. You must fix the random seed so that the input does not wobble when you test resume later.

First count what came in

In /root/runx/runner.py, create scan <작업폴더> (the placeholder stands for the work folder) so that it produces, as JSON, the list of shard names, the total count, and the amount sum. The shard names are in ascending order.

Under in/ in the work folder, pick only the files that end in .jsonl and strip the extension from the name. Do not count blank lines. If the work folder does not exist, you must end with exit code 3 so that the error messages in later steps are honest.

Do not write directly to the destination

Add part <작업폴더> <조각> [--crash=write] (the placeholders stand for the work folder and the shard name). It counts the shard and leaves shard, events, and amount in out/part-<조각>.json, but it writes everything under a temporary name in the same directory as the destination and then swaps it in with os.replace. If you give --crash=write, it dies with exit code 9 just before the swap.

If the temporary name matches the part-*.json listing, that file also goes into the total later. Use a name that starts with a dot. And you must not create the temporary file in /tmp — os.replace fails across filesystems, and that failure is not reproduced on a development machine. The grader looks at whether the destination file is intact after killing it and whether the temporary file remains next to the destination.

Bundle into one run

Add run <작업폴더> --run-id=<이름> (the placeholders stand for the work folder and the run name) to process all the shards in order, rescan the shard outputs, and leave shards, events, and amount in out/total.json. Write total.json with the swap method too.

If you compute the total by adding only the shards processed this time, the shards skipped when resuming later go missing. Always compute the total by rereading all of out/part-*.json. This one rule makes resuming free.

Leave one run as one line

When run finishes, have it append one line to /root/runx/work/ledger.jsonl, and have ledger <작업폴더> (the placeholder stands for the work folder) produce {"runs": 정수, "ok": 정수, "failed": 정수, "last": 마지막 줄} (integers, and the last line).

You only append to the ledger. If you start editing earlier lines, the rule that one run is one line collapses, and at that moment the ledger is no different from a log. Put run_id, started_at, ended_at, status, shards_total, shards_done, events, and amount in a line.

Try killing it at a middle shard

Add --crash-shard=<조각> (the placeholder stands for the shard name) to run. When that shard's turn comes, it leaves a line in the ledger with status failed and failed_shard being that shard, and dies with exit code 9. It does not touch the merged output out/total.json.

You must write the ledger before dying. Without a ledger, all the next person can do is run it again from the start. Put in shards_done only the shards actually finished in this run, and leave total.json untouched — the previous run's answer must remain standing.

Continue from where it died

Add --resume to run. If a shard output out/part-<조각>.json already exists and the shard name inside it is correct, skip that shard, and put the skipped ones in skipped of the response and the ones processed this time in done.

Do not keep a separate mark file; use the shard output itself as the mark. The output is created by a swap, so its existence means that shard certainly finished. If you keep the mark and the output separately, you can end up with only the mark left and the output in a half-written state. The total is still collected again from all the shard outputs.

Keep it from being attached twice

Add commit <작업폴더> --run-id=<이름> (the placeholders stand for the work folder and the run name). Read out/total.json and attach one line of run_id, events, and amount to the day's ledger out/daily.jsonl, but if that run_id is already there, do not attach it and produce {"appended": false, ...}. Commit one run in your own work folder too.

Processing a shard is an overwrite, so it is the same however many times you do it, but attaching one line to the ledger grows every time it is called. If you decide duplicates by time, you cannot tell apart two runs that ran on the same day. Decide by the name of the run.