重跑过去的区间 — 不要算两遍
目标
创建接收区间并重新运行过去的运行器 runner.py。用与定期运行相同的代码以分区为单位替换,用区间预留避免与定期运行冲突,把已经发出的区间封存并以更正的形式留下,并且亲手重现累计汇总因回填而翻倍的情况,再修好它。
为什么重要
回填是把过去的区间重新跑一遍,但难的部分不是计算,而是协调。回填运行的同时,定期运行也在运行。两者同时写同一个日期,就没有人知道哪一边赢了。 同一课程中的幂等实验,是让同一行插入两次结果也相同。这里讲的不是行,而是区间。即使分区表是幂等的,只要累计表不是幂等的,一次回填就会让数字膨胀。累加式的累计遇到回填就一定会出错,而派生式的累计无论运行多少次,得到的值都一样。 而且有些地方不能回头。已经发出的报告和已经送出的告警,都不是数据,而是事件。这些区间要封存,新值出来时,不要覆盖,而要作为更正单独记录,才能同时回答“当时我们发出去的是什么”和“现在什么才是对的”。 评分器不会相信你写出来的文字。它会在临时目录里摆好评分器生成的分区,用环境变量指向那个存储,然后真正运行你的运行器,并直接打开 sqlite 表来核对。日期和金额每次运行都会变化。
步骤
- 创建并运行 /root/backfill/gen_events.py,在 /root/backfill/events 下生成从
dt=2026-02-01.jsonl到dt=2026-02-14.jsonl的十四天分区。 - 在 /root/backfill/runner.py 中实现
init和run <시작일> <끝일> <주인>(占位符依次为开始日期、结束日期、持有者),统计该区间内各分区并写入daily。 - 把
run改成整体替换分区,并把整体运行的答案、按天运行的答案和重新运行的结果写入 /root/backfill/split.json。 - 加上
rollup-add <시작일> <끝일>(占位符依次为开始日期、结束日期)和rollup,并把回填之后两种方式相差多少写入 /root/backfill/double.json。 - 加上
claim <시작일> <끝일> <주인>和release <시작일> <끝일> <주인>,预留区间。 - 让
run跳过被别人占用的分区,并以skipped报告。如果有被跳过的,退出码是 5。 - 加上
seal <시작일> <끝일>和amend <날짜> <사유>(占位符依次为开始日期、结束日期、日期、原因)。被封存的分区,run不会去碰;迟到的单据则以更正的形式留在corrections中。 - 用 /root/backfill/backfill_report.json 和 /root/backfill/backfill_report.md 留下一页报告。
参考
- 运行契约:
python3 /root/backfill/runner.py <명령> ...(占位符为命令)。结果以一个 JSON 对象输出到标准输出。退出码:0 成功,2 命令不认识或参数个数不对,3 没有状态存储,4 预留被拒绝,5 有被跳过的区间,6 对未封存的区间调用了amend。 - 状态存储路径必须可以通过环境变量
BACKFILL_DB修改,原始分区目录必须可以通过BACKFILL_EVENTS修改。默认值分别是/root/backfill/state.db和/root/backfill/events。评分器会用这两个变量指向它自己的存储。 - 分区文件名是
dt=YYYY-MM-DD.jsonl,行的格式是{"order_id": 문자열, "dt": 날짜, "amount": 정수, "status": "paid" 또는 "refund"}(占位符依次为字符串、日期、整数、表示“或”的词)。退款行的 amount 是负数。一个分区的orders是行数,amount是 amount 的合计。 - 一共有五张表。
daily(dt, orders, amount, owner)、rollup(metric, value)、claims(dt, owner)、seals(dt)、corrections(dt, orders_before, amount_before, orders_after, amount_after, reason)。daily不设主键——亲手实现替换,正是本实验的要点。 run的响应是{"owner": 문자열, "done": [날짜...], "skipped": [[날짜, 사유]...], "orders": 정수, "amount": 정수}(占位符依次为字符串、日期、日期和原因、整数、整数)。orders和amount是只把done中的分区相加得到的值。原因有claimed_by:<주인>(占位符为持有者)、sealed、no_data三种。原始文件不存在的日期(no_data)不会让退出码变成 5。claim的响应是{"owner": 문자열, "claimed": [날짜...], "denied": [[날짜, 주인]...]}(占位符依次为字符串、日期、日期和持有者)。自己已经占用的日期会原样放进claimed。release的响应是{"owner": 문자열, "released": [날짜...]}(占位符依次为字符串、日期),只释放自己的预留。seal的响应是{"sealed": [날짜...]}(占位符为日期)。amend的响应是{"dt": 날짜, "orders_before": 정수, "amount_before": 정수, "orders_after": 정수, "amount_after": 정수, "reason": 문자열}(占位符依次为日期、整数、整数、整数、整数、字符串)。amend不修改daily,只往corrections里加一行。rollup-add的响应和rollup的响应是{"mode": "add" 또는 "derive", "order_total": 정수, "amount_total": 정수}(占位符依次为表示“或”的词、整数、整数)。rollup-add把区间合计加到已有值上,rollup则从daily全部数据重新计算并覆盖。- 不按性能来判定。不要测量回填花了多长时间,只留下动了什么、按什么顺序动的。
- 官方文档:python sqlite3 · SQLite Transaction · SQLite UPSERT · Airflow Dag Runs
- 常见错误:另外写回填用的代码、把整个区间放进一个事务、悄悄覆盖别人的区间、把累计保持为累加的方式、覆盖已经发出的区间而丢掉当时发出了什么。
- 想直接看状态存储,可以使用
sqlite3 /root/backfill/state.db 'SELECT * FROM daily ORDER BY dt'。
生成十四天的分区
创建并运行 /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 步测得的值。报告里要用数字写出相差的金额。