凌晨三点挂掉的结算作业 — 台账与原子替换
目标
创建一个即使中途终止也安全的运行器 runner.py。先用临时名字写产出、再改名替换,避免留下写到一半的文件;把每次运行作为一行记入账本;从终止的分片开始续跑;并且让同一次运行提交两次,数字也不会增加。
为什么重要
管道一定会在中途挂掉。问题不在于它会挂掉,而在于挂掉之后会留下什么。如果一直是直接写目标文件,就会留下写到一半的文件,这个文件的大小和名字都完好无损,下一步会照常把它读走。 重新运行同样危险。如果没有记录说明上一次运行做到了哪里,就只能从头再跑一遍,而在最后追加到账簿的地方,同一笔金额会被加两次。这就是“明明重新跑了一遍,为什么变成了两倍”的真相。 本实验构建的机制有三个。第一个是原子替换:在同一目录下用临时名字写完,再用 os.replace 改名替换。第二个是运行账本:每次运行追加一行,如果失败了,就记下死在哪个分片。第三个是续跑:把已完成分片的产出本身当作标记,直接跳过。 本课程的 dp-idempotent 实验讲的是数据库一侧的幂等性——即同一行插入两次也只会得到一行。本实验是它的前一步:讲的是进程终止的地方,文件系统上会留下什么,以及看着留下的东西,应该从哪里重新开始。 评分器不会相信你写出来的文字。它会在临时目录里摆好评分器自己生成的输入分片,然后真正运行你的运行器。故意把它杀掉之后,还会检查目标文件是否原封不动、临时文件是否留在了目标文件旁边、账本里是否写下了失败的分片名。分片数和金额每次运行都会变化。
步骤
- 创建并运行 /root/runx/gen_shards.py,在 /root/runx/work/in 下生成输入分片。
- 在 /root/runx/runner.py 中实现
scan,输出分片列表、记录数和合计。 - 加上
part,处理一个分片,并让产出先用临时名字写入、再改名替换。另外用--crash=write做出一条在改名替换前一刻终止的路径。 - 加上
run,处理全部分片并生成合并后的产出。 - 每次运行都在账本中留下一行,并用
ledger输出汇总。 - 用
--crash-shard让它在中间的分片处终止,并确认账本里留下了失败的分片、合并后的产出没有被动过。 - 加上
--resume,跳过已完成的分片,从终止的位置接着运行。 - 加上
commit,让同一次运行不会被两次追加到当天的账簿 /root/runx/work/out/daily.jsonl 中。
参考
- 所有工作都在
/root/runx下进行。工作目录是/root/runx/work。 - 工作目录结构:输入是
in/<조각>.jsonl(占位符为分片名),分片产出是out/part-<조각>.json,合并后的产出是out/total.json,当天的账簿是out/daily.jsonl,账本是ledger.jsonl。分片名就是输入文件名去掉.jsonl之后的部分。 - 输入的每一行是一个 JSON 对象,
amount字段里是整数。也可以有其他字段。 - 运行契约:
python3 /root/runx/runner.py <명령> <작업폴더> [...](占位符依次为命令、工作目录)。结果以一个 JSON 对象输出到标准输出。成功时退出码为 0,缺少工作目录或所需文件时为 3,用法错误时为 2,进入故意终止的路径时为 9。 scan <작업폴더>的响应是{"shards": [이름 오름차순], "events": 정수, "amount": 정수}(占位符依次为工作目录;按名称升序排列、整数、整数)。part <작업폴더> <조각> [--crash=write]的响应是{"shard": 이름, "events": 정수, "amount": 정수, "path": 산출물 경로}(占位符依次为工作目录、分片名;名称、整数、整数、产出路径)。产出文件里包含 shard、events、amount。run <작업폴더> --run-id=<이름> [--resume] [--crash-shard=<조각>]的响应是{"run_id": 이름, "status": "ok", "shards_total": 정수, "shards_done": [이름], "skipped": [이름], "done": [이름], "events": 정수, "amount": 정수, "started_at": 문자열, "ended_at": 문자열}(占位符依次为工作目录、运行名、分片名;响应中依次为运行名、整数、分片名列表、分片名列表、分片名列表、整数、整数、字符串、字符串)。skipped是续跑时被跳过的分片,done是这一次处理的分片。- 账本的一行包含 run_id、started_at、ended_at、status、shards_total、shards_done、events、amount,如果失败,再加上 failed_shard。status 要么是
ok,要么是failed。 ledger <작업폴더>的响应是{"runs": 정수, "ok": 정수, "failed": 정수, "last": 마지막 원장 줄}(占位符依次为工作目录;整数、整数、整数、账本的最后一行)。commit <작업폴더> --run-id=<이름>的响应是{"appended": 참거짓, "run_id": 이름, "lines": 장부 줄 수}(占位符依次为工作目录、运行名;布尔值、运行名、账簿行数)。账簿的一行是 run_id、events、amount。- 临时文件要建在与目标相同的目录里,并且名字不能被
part-*.json列表匹配到。os.replace 跨文件系统会失败。 - 官方文档:os.replace · rename(2) · SQLite Atomic Commit · python json
- 常见错误:直接打开目标文件写入、把临时文件建在
/tmp里、只累加这一次运行处理的部分来得出总计、按时间判断账簿是否重复。 - 想看终止之后留下了什么,可以使用
ls -a /root/runx/work/out。
创建上游丢下来的分片
创建并运行 /root/runx/gen_shards.py,在 /root/runx/work/in 下生成分片文件。至少要有 4 个分片,每个分片至少 5 行,总共至少 40 行,每一行是包含 id 和整数 amount 的一个 JSON 对象。
一个分片就是一个 JSON Lines 文件。文件名去掉 .jsonl 之后就是分片名,所以要像 h00、h01 这样补齐位数,使排序后就是时间顺序。要固定随机数种子,这样之后测试续跑时,输入才不会变动。
先数一数进来了什么
在 /root/runx/runner.py 中实现 scan <작업폴더>(占位符为工作目录),以 JSON 输出分片名列表、总记录数和金额合计。分片名按升序排列。
只挑选工作目录下 in/ 里以 .jsonl 结尾的文件,并去掉文件名的扩展名。空行不计入。如果工作目录不存在,就要以退出码 3 结束,这样后面步骤的错误信息才能如实反映问题。
不直接写目标文件
加上 part <작업폴더> <조각> [--crash=write](占位符依次为工作目录、分片名)。统计分片后,把 shard、events、amount 写入 out/part-<조각>.json,但要在与目标相同的目录里用临时名字写完,再用 os.replace 改名替换。传入 --crash=write 时,在改名替换的前一刻以退出码 9 终止。
如果临时名字被 part-*.json 列表匹配到,后面汇总时连这个文件也会被算进合计。使用以点号开头的名字。另外不能把临时文件建在 /tmp 里——os.replace 跨文件系统会失败,而这个失败在开发机上是无法复现的。评分器会在终止之后检查目标文件是否完好无损、临时文件是否留在了目标文件旁边。
合并成一次运行
加上 run <작업폴더> --run-id=<이름>(占位符依次为工作目录、运行名),按顺序处理全部分片,再重新遍历分片产出,把 shards、events、amount 写入 out/total.json。total.json 也用改名替换的方式写入。
如果只累加这一次处理的分片来得出总计,之后续跑时被跳过的分片就会漏掉。总计要始终重新读取 out/part-*.json 全部文件来得出。仅这一条规则,就能让续跑几乎不费力气。
把一次运行记成一行
让 run 结束时在 /root/runx/work/ledger.jsonl 中追加一行,并用 ledger <작업폴더> 输出 {"runs": 정수, "ok": 정수, "failed": 정수, "last": 마지막 줄}(占位符依次为工作目录;整数、整数、整数、最后一行)。
账本只追加。一旦开始修改前面的行,“一次运行一行”的规则就被破坏,到那时账本和日志就没有区别了。在一行里包含 run_id、started_at、ended_at、status、shards_total、shards_done、events、amount。
在中间的分片处终止看看
给 run 加上 --crash-shard=<조각>(占位符为分片名)。轮到那个分片时,在账本中留下 status 为 failed、failed_shard 为该分片的一行,然后以退出码 9 终止。不要动合并后的产出 out/total.json。
必须先写账本,再终止。没有账本的话,下一个人能做的就只有从头重新运行。shards_done 里只放这一次运行中真正完成的分片,而 total.json 保持原样不动——上一次运行的答案必须原封不动地留在那里。
从终止的位置接着运行
给 run 加上 --resume。如果分片产出 out/part-<조각>.json(占位符为分片名)已经存在,并且其中的 shard 名称正确,就跳过该分片,被跳过的放入响应的 skipped,这一次处理的放入 done。
不要另设标记文件,而是把分片产出本身当作标记。产出是通过改名替换生成的,所以它存在就意味着该分片一定已经完成。如果把标记和产出分开放,就可能出现只剩标记、产出却处于写到一半状态的情况。总计仍然是从全部分片产出中重新汇总。
避免被追加两次
加上 commit <작업폴더> --run-id=<이름>(占位符依次为工作目录、运行名)。读取 out/total.json,把包含 run_id、events、amount 的一行追加到当天的账簿 out/daily.jsonl;如果该 run_id 已经存在,就不追加,而是输出 {"appended": false, ...}。在自己的工作目录里也提交一次运行。
分片处理是覆盖写入,做多少次结果都一样,而往账簿里追加一行,每调用一次就会多一行。如果按时间判断是否重复,就区分不了同一天运行了两次的情况。按运行的名字来判断。