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

数据流水线

重跑过去的区间 — 不要算两遍

在 TT Lab 中继续学习

目标

创建接收区间并重新运行过去的运行器 runner.py。用与定期运行相同的代码以分区为单位替换,用区间预留避免与定期运行冲突,把已经发出的区间封存并以更正的形式留下,并且亲手重现累计汇总因回填而翻倍的情况,再修好它。

为什么重要

回填是把过去的区间重新跑一遍,但难的部分不是计算,而是协调。回填运行的同时,定期运行也在运行。两者同时写同一个日期,就没有人知道哪一边赢了。 同一课程中的幂等实验,是让同一行插入两次结果也相同。这里讲的不是行,而是区间。即使分区表是幂等的,只要累计表不是幂等的,一次回填就会让数字膨胀。累加式的累计遇到回填就一定会出错,而派生式的累计无论运行多少次,得到的值都一样。 而且有些地方不能回头。已经发出的报告和已经送出的告警,都不是数据,而是事件。这些区间要封存,新值出来时,不要覆盖,而要作为更正单独记录,才能同时回答“当时我们发出去的是什么”和“现在什么才是对的”。 评分器不会相信你写出来的文字。它会在临时目录里摆好评分器生成的分区,用环境变量指向那个存储,然后真正运行你的运行器,并直接打开 sqlite 表来核对。日期和金额每次运行都会变化。

步骤

  1. 创建并运行 /root/backfill/gen_events.py,在 /root/backfill/events 下生成从 dt=2026-02-01.jsonl 到 dt=2026-02-14.jsonl 的十四天分区。
  2. 在 /root/backfill/runner.py 中实现 init 和 run <시작일> <끝일> <주인>(占位符依次为开始日期、结束日期、持有者),统计该区间内各分区并写入 daily。
  3. 把 run 改成整体替换分区,并把整体运行的答案、按天运行的答案和重新运行的结果写入 /root/backfill/split.json。
  4. 加上 rollup-add <시작일> <끝일>(占位符依次为开始日期、结束日期)和 rollup,并把回填之后两种方式相差多少写入 /root/backfill/double.json。
  5. 加上 claim <시작일> <끝일> <주인> 和 release <시작일> <끝일> <주인>,预留区间。
  6. 让 run 跳过被别人占用的分区,并以 skipped 报告。如果有被跳过的,退出码是 5。
  7. 加上 seal <시작일> <끝일> 和 amend <날짜> <사유>(占位符依次为开始日期、结束日期、日期、原因)。被封存的分区,run 不会去碰;迟到的单据则以更正的形式留在 corrections 中。
  8. 用 /root/backfill/backfill_report.json 和 /root/backfill/backfill_report.md 留下一页报告。

参考

生成十四天的分区

创建并运行 /root/backfill/gen_events.py,在 /root/backfill/events 下生成从 dt=2026-02-01.jsonl 到 dt=2026-02-14.jsonl 的文件。一天一个文件,退款行的 amount 是负数。

日期可以用 datetime.date.fromisoformat 和 timedelta 递增。如果让每一天的行数不同,之后各分区的合计就能互相区分。掺进几条退款行,让金额不是单纯的累计。

创建接收区间并运行的运行器

在 /root/backfill/runner.py 中实现 init 和 run <시작일> <끝일> <주인>(占位符依次为开始日期、结束日期、持有者)。run 统计该区间内每个分区的行数和金额并写入 daily,对于原始文件不存在的日期,在 skipped 中以 no_data 记录。

定期运行只不过是区间为一天的情况。所以只写一个函数。表在 init 里用 CREATE TABLE IF NOT EXISTS 创建,daily 不设主键。no_data 不会让退出码变成 5。

以分区为单位替换

把 run 改成在重新运行一个分区时,先删除该分区的旧结果再写入。然后把整体运行的答案、按天运行的答案和重新运行的结果,以 whole、by_day、rerun_changed 写入 /root/backfill/split.json,并把 /root/backfill/state.db 的 daily 也重新填充成每个日期一行。

替换就是删除该日期的行再重新写入。要比较整体运行的答案和按天运行的答案,需要把存储新建两次,而用 BACKFILL_DB 环境变量指向临时存储,就可以在不动正式存储的情况下进行比较。

重现累计汇总翻倍

加上 rollup-add <시작일> <끝일>(占位符依次为开始日期、结束日期)和 rollup。然后在临时存储里运行全部区间并加到累计里,回填一部分区间之后再把该区间加一遍,把它与派生出的值相差多少,以 backfill_range、true_total、add_after_backfill、gap 写入 /root/backfill/double.json。

rollup-add 把区间合计加到已有值上,rollup 则把 daily 全部重新统计并覆盖。相差的幅度必须恰好等于回填区间的金额合计——因为那个区间被计入了累计两次。

先预留区间

加上 claim <시작일> <끝일> <주인>(占位符依次为开始日期、结束日期、持有者)和 release <시작일> <끝일> <주인>。被别人占用的日期连同持有者一起放进 denied,退出码是 4。release 只释放自己的预留。

如果把 claims 表的主键设为日期,同一个日期就不可能被两个持有者占用。自己已经占用的日期再次占用,不要当作拒绝,而要当作成功——重试如果看起来是失败,就没有人会去重试。

遇到别人的区间就跳过并报告

让 run 不去碰被别人占用的分区,并在 skipped 中以 claimed_by:<주인>(占位符为持有者)记录。如果有被跳过的,退出码是 5。自己占用的分区照常运行。

要点是既不等待也不覆盖。跳过并说出来,调用方就能决定是再次调用还是找人处理。预留列表只需要在遍历区间之前读取一次。

已经发出的区间以更正的形式留下

加上 seal <시작일> <끝일> 和 amend <날짜> <사유>(占位符依次为开始日期、结束日期、日期、原因)。被封存的分区,run 会以 sealed 跳过。然后在 2026-02-02 分区里加一行迟到的单据,把这个日期封存,再用 amend 留下更正。更正的 before 和 after 必须不同。

amend 不修改 daily。它只是把从原始数据重新统计出的值和 daily 中留存的值并排写进 corrections。如果对没有封存的日期调用 amend,要以退出码 6 拒绝——这样的日期直接 run 就行了。

把重新运行的这一轮写成一页报告

在 /root/backfill/backfill_report.json 中写入 partitions、orders、amount、sealed、corrections、rollup、double_gap,并在 /root/backfill/backfill_report.md 中写成四节,标题是 ## 무엇을 다시 돌렸나(韩文,意为“重新运行了什么”)、## 정기 실행과 어떻게 부딪혔나(韩文,意为“和定期运行发生了怎样的冲突”)、## 누적 집계는 왜 두 배가 되나(韩文,意为“累计汇总为什么会翻倍”)、## 되돌리면 안 되는 자리(韩文,意为“不能回头的地方”)。

partitions、orders、amount 从 daily 表中读取,sealed 从 seals 中读取,corrections 取行数,rollup 从派生出的累计中读取。double_gap 直接使用第 4 步测得的值。报告里要用数字写出相差的金额。