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

数据流水线

重跑过去 — 回填把数字算成两倍的时候

在 TT Lab 中继续学习

一句话总结

回填不是“把过去的区间重新跑一遍”,而是“用和定期运行相同的代码,不对同一个分区动两次手,并且不覆盖已经发出去的数字,重新跑一遍”。

为什么需要它

出了个 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만원   ← 몇 번 돌려도 같다

在 2 月 1 日到 14 日的日期轴上,定期运行的区间和 9 日到 11 日的回填区间重叠的示意图。如果采用累加方式的累计,5,182 万韩元再加上 842 万韩元就变成 6,024 万韩元,842 万韩元被算了两次;如果采用重新统计的累计,仍然是 5,182 万韩元,无论运行多少次都一样

即使分区表是幂等的,累计表不是幂等的也没有用。累计不要累加,而要派生。从原始分区重新统计,无论回填运行多少次,值都不会摇摆。

第五条是:有些地方不能回头。已经发出的报告、已经送出的告警、已经结算的金额,都不是数据,而是事件。这些区间要封存,新值出来时,不要覆盖,而要作为更正单独记录。这样才能同时回答“当时我们发出去的是什么”和“现在什么才是对的”。

在现场相遇的样子

第一,另外做一个回填脚本。为了着急时用一次而写的脚本留了下来,半年后还在运行。出现了只有那个脚本才知道的异常处理,两条路径的答案就分叉了。

第二,把整个区间放进一个事务。14 天的数据一次性提交的话,在第 13 天失败就得从头再来,而在此期间状态存储对此一无所知。如果每个分区各自完成,重新开始的位置就白白得到了。

第三,不做预留就相信“现在不是定期运行的时间”。重试、手动运行、时间差,都会让这种相信落空。预留是由代码遵守的约定,而时间段是由人遵守的约定。

第四,只把回填范围记在日志里。如果状态存储里没有留下什么时候是谁把什么重新跑了一遍,数字不对的时候就没有可以回溯的依据。而且不要用运行时间来推断是不是回填——即使在同一台机器上,运行时间也会波动到两倍。

实际工作中真正重要的事

下一项实验要做什么

生成十四天的订单分区,再一步步扩充运行器 runner.py。让它接收区间并以分区为单位替换,确认整体运行得到的答案和按天运行得到的答案相同,并用数字记下累加式累计和派生式累计在回填之后相差多少。然后加上区间预留,让定期运行跳过别人的区间并报告;最后把已经发出的区间封存,对于迟到的单据,不是覆盖,而是以更正的形式留下。评分器每次都会用不同的日期和金额生成自己的分区,真正运行你的运行器,并直接读取状态存储来核对。