TT Lab
开始
学习 学习路径 课程

数据流水线

凌晨三点挂掉的结算作业 — 台账与原子替换

在 TT Lab 中继续学习

目标

创建一个即使中途终止也安全的运行器 runner.py。先用临时名字写产出、再改名替换,避免留下写到一半的文件;把每次运行作为一行记入账本;从终止的分片开始续跑;并且让同一次运行提交两次,数字也不会增加。

为什么重要

管道一定会在中途挂掉。问题不在于它会挂掉,而在于挂掉之后会留下什么。如果一直是直接写目标文件,就会留下写到一半的文件,这个文件的大小和名字都完好无损,下一步会照常把它读走。 重新运行同样危险。如果没有记录说明上一次运行做到了哪里,就只能从头再跑一遍,而在最后追加到账簿的地方,同一笔金额会被加两次。这就是“明明重新跑了一遍,为什么变成了两倍”的真相。 本实验构建的机制有三个。第一个是原子替换:在同一目录下用临时名字写完,再用 os.replace 改名替换。第二个是运行账本:每次运行追加一行,如果失败了,就记下死在哪个分片。第三个是续跑:把已完成分片的产出本身当作标记,直接跳过。 本课程的 dp-idempotent 实验讲的是数据库一侧的幂等性——即同一行插入两次也只会得到一行。本实验是它的前一步:讲的是进程终止的地方,文件系统上会留下什么,以及看着留下的东西,应该从哪里重新开始。 评分器不会相信你写出来的文字。它会在临时目录里摆好评分器自己生成的输入分片,然后真正运行你的运行器。故意把它杀掉之后,还会检查目标文件是否原封不动、临时文件是否留在了目标文件旁边、账本里是否写下了失败的分片名。分片数和金额每次运行都会变化。

步骤

  1. 创建并运行 /root/runx/gen_shards.py,在 /root/runx/work/in 下生成输入分片。
  2. 在 /root/runx/runner.py 中实现 scan,输出分片列表、记录数和合计。
  3. 加上 part,处理一个分片,并让产出先用临时名字写入、再改名替换。另外用 --crash=write 做出一条在改名替换前一刻终止的路径。
  4. 加上 run,处理全部分片并生成合并后的产出。
  5. 每次运行都在账本中留下一行,并用 ledger 输出汇总。
  6. 用 --crash-shard 让它在中间的分片处终止,并确认账本里留下了失败的分片、合并后的产出没有被动过。
  7. 加上 --resume,跳过已完成的分片,从终止的位置接着运行。
  8. 加上 commit,让同一次运行不会被两次追加到当天的账簿 /root/runx/work/out/daily.jsonl 中。

参考

创建上游丢下来的分片

创建并运行 /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, ...}。在自己的工作目录里也提交一次运行。

分片处理是覆盖写入,做多少次结果都一样,而往账簿里追加一行,每调用一次就会多一行。如果按时间判断是否重复,就区分不了同一天运行了两次的情况。按运行的名字来判断。