It Dies Halfway: Run Ledgers and Atomic Swaps
One-line summary
A pipeline certainly dies midway. What remains at that point, a half-written output or a record of how far it got, decides the next morning.
Why this was needed
A settlement job died at 3 a.m. When you come to work in the morning, the result file is there. Its size looks plausible. But the numbers are strange. If you open the file, the last line is cut off in the middle — the process vanished while writing.
Here the second incident happens. You think "I'll just rerun it" and rerun. And that day's sales are counted twice. Nobody knows how far the previous run got, and the rerun goes through everything from start to finish.
The third incident is quieter. A temporary file left by the dead run is picked up together in the next step's listing, and the same shard is added two more times. Nobody sees an error.
These three do not arise because the code is wrong. They arise from code written thinking only of the success case.
How it works
There are three devices that prevent this.
First, atomic replacement. Do not write directly to the destination file. Write everything under a temporary name in the same directory, and then rename it into the destination. In Python that is os.replace. What it calls underneath is rename(2), and that documentation states that if the new name already exists it is replaced atomically, so that "there is no moment when the name does not exist when another process looks for it." So the reader, whenever it looks, sees either the whole old file or the whole new file. It cannot see a half-written state.
There is a condition, though. The temporary file and the destination must be in the same filesystem. The same documentation notes that it fails if the two paths are on different mounts (EXDEV). This is why code that writes to /tmp and moves to the result directory runs fine on a development machine and then dies in production. So you create the temporary file right next to the destination.
There is a rule for the temporary name too. If the reader builds its listing with part-*.json, the temporary name must not match it. Start it with a dot or attach a different suffix.
Second, the run ledger. Write one run as one line. When it started, up to which shard it finished, whether it succeeded or failed, and if it failed, where it died. You only append and never edit. That is how you can trace "yesterday's run" later. What SQLite's atomic commit documentation explains is ultimately the same structure — keep the in-progress state separately and change the mark all at once when it is done.
Third, per-shard marks. To avoid starting from the beginning when you run again, there has to be a record of how far you got. You can have a separate mark file, or use each shard's output itself as the mark. The latter is safer — because the output is replaced atomically, the existence of the output means that shard certainly finished. If the mark and the output are separate, you can end up with only the mark left and the output in a half-written state.
나쁜 순서 좋은 순서
결과 파일을 열고 쓴다 임시 파일에 다 쓴다
... 여기서 죽음 ... 여기서 죽어도 목적지는 멀쩡
닫는다 os.replace 로 바꿔 단다
=> 반쯤 쓰인 파일이 남는다 => 옛 파일 아니면 새 파일
What it looks like in the field
First, resume and rerun are different words. A rerun is running again from the start, and a resume is skipping the shards that are finished. To support resume, you must be able to read "this shard is finished." Writing it in a log is not enough. A log is for people to read, and a resume is decided by a program.
Second, you collect the total again. If you compute the total by adding only what the resumed run processed this time, the skipped shards are missing. You always compute the total by rescanning all the shard outputs. That way a resumed run and a run that finished in one go give the same answer.
Third, the place where things get added twice is usually the last one. Processing a shard is an overwrite, so it is the same however many times you do it, but attaching one line to the day's ledger is an append, so it grows every time it is called. So before attaching, you check whether this run has already been attached. The basis for the decision is not the time but the name of the run. If you decide by time, you cannot tell apart two runs that ran on the same day.
Fourth, without knowing where it died, you can do nothing. If the ledger has no name of the failed shard, all the next person can do is run it again from the start. For a six-hour job, that difference is large.
What really matters in practice
- Do not write directly to the destination. Write under a temporary name and swap it in within the same filesystem.
- Keep the temporary name from matching the reader's listing.
- Give each run a name. Tell duplicates apart by name, not by time.
- Write the failed shard in the ledger. That one cell makes resume possible.
- Collect the total again from the outputs. Do not add only what this run did.
What to do in the next lab
You turn the nightly settlement input into shards and build up the runner runner.py step by step. You have it write shard outputs under a temporary name and swap them in, deliberately kill it midway through writing, and confirm that the destination is intact. You attach a run ledger, kill it at a middle shard to see that the failed shard remains in the ledger, and resume from that point. Finally, you make it so that committing the same run twice does not grow the day's ledger. The grader sets up its own input with a different number of shards and amounts each time, actually runs your runner, and even checks what remains after killing it.