重跑过去 — 回填把数字算成两倍的时候
一句话总结
回填不是“把过去的区间重新跑一遍”,而是“用和定期运行相同的代码,不对同一个分区动两次手,并且不覆盖已经发出去的数字,重新跑一遍”。
为什么需要它
出了个 bug。过去两周的汇总是错的。用修好的代码把那个区间重新跑一遍。到这里为止,谁都会做。
问题出在后面。回填运行的同时,定期运行也在继续。两者同时写同一个日期,就没有人知道哪一边赢了。累计汇总会随着回填加进去的量而膨胀。而且其中三天的数据已经作为报告发给了客户。
这三件事是各自独立的问题,要用各自不同的机制来防范。本课程前面实验讲的行级幂等在这里只是必要条件,而不是充分条件。即使让同一行插入两次结果也相同,区间级别的协调仍然必须另外具备。
工作原理
第一条原则是:回填必须和定期运行是同一份代码。如果另外放一个只用于回填的脚本,两份代码就会慢慢分叉,到了某一天,没有人能解释回填结果为什么和定期运行的结果不同。所以运行器只要会做一件事就够了:“接收一个区间,重新计算该区间的分区”。定期运行只不过是这个区间为昨天一天的情况。Airflow 的回填也是针对过去的区间运行同一个 DAG,而不是去调用另一个 DAG。
第二条是:以分区为单位整体替换。把一个分区的结果删掉再重新写入,那么无论是整个区间一起跑,还是按天拆开跑,得到的都是同一个答案。没有这个性质,回填就无法中途停下再重新开始。而且把分区拆开跑通常更好——失败的时候,做到了哪里会以分区边界的形式显露出来。
第三条是:先预留区间。回填会用到哪些日期,先写进状态存储;定期运行则不去碰被别人占用的日期,跳过它们,并把这件事报告出来。与其安静地等待或安静地覆盖,跳过并说出来要好。在使用锁的时候,像 SQLite 的 BEGIN IMMEDIATE 那样从一开始就表明写入意图,比晚一步失败要好,道理也是一样的。
第四条是:不要把回填混进累计汇总里。这里是最常出错的地方。
누적을 '더하는' 방식
1) 2월 1일부터 14일까지 정기 실행 누적 += 5,182만원 → 5,182만원
2) 2월 9일부터 11일까지 백필 daily 는 제자리에 치환됨
3) 백필 구간을 누적에 더함 누적 += 842만원 → 6,024만원 ← 842만원이 두 번
누적을 '다시 세는' 방식
daily 표 전체를 합쳐 넣는다 누적 = 5,182만원 ← 몇 번 돌려도 같다
即使分区表是幂等的,累计表不是幂等的也没有用。累计不要累加,而要派生。从原始分区重新统计,无论回填运行多少次,值都不会摇摆。
第五条是:有些地方不能回头。已经发出的报告、已经送出的告警、已经结算的金额,都不是数据,而是事件。这些区间要封存,新值出来时,不要覆盖,而要作为更正单独记录。这样才能同时回答“当时我们发出去的是什么”和“现在什么才是对的”。
在现场相遇的样子
第一,另外做一个回填脚本。为了着急时用一次而写的脚本留了下来,半年后还在运行。出现了只有那个脚本才知道的异常处理,两条路径的答案就分叉了。
第二,把整个区间放进一个事务。14 天的数据一次性提交的话,在第 13 天失败就得从头再来,而在此期间状态存储对此一无所知。如果每个分区各自完成,重新开始的位置就白白得到了。
第三,不做预留就相信“现在不是定期运行的时间”。重试、手动运行、时间差,都会让这种相信落空。预留是由代码遵守的约定,而时间段是由人遵守的约定。
第四,只把回填范围记在日志里。如果状态存储里没有留下什么时候是谁把什么重新跑了一遍,数字不对的时候就没有可以回溯的依据。而且不要用运行时间来推断是不是回填——即使在同一台机器上,运行时间也会波动到两倍。
实际工作中真正重要的事
- 回填和定期运行是同一个函数。不同的只能是区间参数。
- 以分区为单位替换,每个分区各自提交。这样就有了重新开始的位置。
- 预留区间,遇到别人的区间就跳过并报告。不等待,也不覆盖。
- 累计不要累加,而要派生。累加式的累计遇到回填就一定会出错。
- 已经发出去的要封存,并以更正的形式留下。一旦覆盖,当时的事实就消失了。
下一项实验要做什么
生成十四天的订单分区,再一步步扩充运行器 runner.py。让它接收区间并以分区为单位替换,确认整体运行得到的答案和按天运行得到的答案相同,并用数字记下累加式累计和派生式累计在回填之后相差多少。然后加上区间预留,让定期运行跳过别人的区间并报告;最后把已经发出的区间封存,对于迟到的单据,不是覆盖,而是以更正的形式留下。评分器每次都会用不同的日期和金额生成自己的分区,真正运行你的运行器,并直接读取状态存储来核对。