Harness 工程进展:合并、事件与恢复 · 阅读系列导读 →
一个观察扫描器定期读取信号文件。同一种情况反复出现时,它会累加对应的计数,供后续判断使用。文件里已有两条属于同一个键 K 的信号,扫描器把它们读进内存,准备将 K 的累计值从 0 改成 2。就在这几行程序之间,进程退出了。
重启以后,系统究竟该从哪里继续?如果它认为两条信号已经处理,却没有保存计数,这次观察就消失了;如果它保存了计数,又把两条信号当作新输入,累计值便会变成 4。两种错误都可能发生在没有网络、只有本地文件的程序里。决定结果的不是“有无重试”,而是重启后留下了什么证据。
这篇文章讨论 Harness 的观察信号计数路径。它面对的是一个固定、只向末尾追加的信号文件,以及可能在处理途中退出的本地进程。我们沿着这个案例推导写入顺序:读取时保留旧游标,将累计值和已计字节水位一起保存,最后推进游标。这个顺序允许信号再次被读到,并让已经完成的计数在重放时被识别出来。
两份状态,回答两个问题
扫描器需要保存读取游标 C。它是信号文件中的字节偏移,表示下次从哪个位置继续读。追加内容不会改变既有字节的位置,因此游标无需每轮重新计算之前的行数。它也需要保存计数表 N,例如 N[K] = 2,表示这个键已经累计了两次呈现。
这两份状态最容易被混成一句“已经处理到这里”。其实,C 回答的是读取位置,N 回答的是已经产生的计数结果。文件读到了末尾,只能说明数据进入了本轮的内存;内存里的两条记录和新累计值,会随进程一起消失。恢复必须依据保存在文件里的状态,不能依据上一轮执行到哪一行的印象。
最直接的实现是读完后立刻把 C 推到文件末尾,再更新 N。正常执行时,它看起来很顺畅。但若计数保存遇到 OSError,或者进程在两步之间退出,磁盘上就会留下“游标已经越过两条信号,计数仍为 0”的组合。下一轮从末尾开始,两条信号没有机会再进入计数路径。
把顺序倒过来,先保存 N,再推进 C,就消除了这个丢信号窗口。可是新窗口随之出现:N 已经是 2,C 仍停在 0。重启会再次读出原来的两条信号,如果仍执行两次加一,结果就是 4。确认读取之前保存结果,解决了遗漏;它还需要一个能辨认重放的状态。
| 写入安排与故障位置 | 重启时留下的状态 | 下一轮的后果 |
|---|---|---|
| 先推进游标,计数保存前失败 | C 已前移,N 仍旧 | 信号不再被读到,计数遗漏 |
| 先保存计数,游标推进前失败 | N 已增加,C 仍旧 | 信号重读,直接累加会重复计数 |
| 计数与水位一起保存,游标推进前失败 | N 已增加,水位已前移,C 仍旧 | 信号重读,旧信号跳过计数 |
表里的第三行保留了第二行的重读行为,却改变了重读后的判断。要做到这一点,新增的证据必须与计数结果有相同的保存边界。
从故障窗口推导保存顺序
这份证据叫已计水位 A:它记录已经纳入计数的信号,最远到达哪个行末字节偏移。C 与 A 都以字节为单位,但它们承担不同职责。C 决定读哪些字节;A 决定读出来的普通信号还需不需要加到 N 里。一次故障可以使 C 落后于 A,恢复逻辑正要处理这个差距。
假设把 A 另存为一个文件,先保存 N、后保存 A。进程仍可能在中间退出,留下新计数与旧水位,重放又会双计。若先保存 A、后保存 N,中间退出则留下新水位与旧计数,重放会跳过尚未计入的信号。原先两个文件之间的矛盾,只是换了一对文件。
因此,计数表和已计水位必须组成同一个保存单位。当前实现将它们放进同一份 JSON 状态,先把完整的新内容写到同目录的临时文件,再用 os.replace 替换目标文件。Python 对成功的替换规定了原子操作语义。对这里的进程退出问题,恢复读取的是替换前或替换后的整份状态,计数与水位无需分别猜测。Python 官方文档:os.replace 这里针对进程退出恢复;当前写入未调用 fsync,掉电持久性不在这个保证内。
整份替换也不能独自处理并发更新。两个写者若同时读取旧值 0,各自在自己的副本里加一,再依次替换文件,最后仍可能只有 1。实现用同一把专用 flock 锁包住读取旧状态、修改计数与水位、替换文件这一整段过程。锁约束参与该约定的写者之间的读改写顺序,原子替换约束一次状态更新的可见边界。
图 1 把这两个保存边界放在同一条时间轴上。按这个顺序,进程在两次保存之间退出,会留下“计数状态已更新,读取游标未更新”的组合,水位足以解释它;“读取游标已更新,计数状态未更新”则不应由正常路径产生。
- 01
读入内存,保留旧游标
C=0 · N[K]=0 · A=0 - 02
先保存计数与水位
N[K]=2 · A=80flock包住完整读改写。 - ×
进程在此退出
下轮重读,旧信号不再加一。C=0 · N[K]=2 · A=80 - 03
最后保存读取游标
C=80 · N[K]=2 · A=80
下面的伪代码只展示有效普通信号的计数分支,保留决定恢复行为的顺序。信号列表在读取阶段已按文件顺序得到,每条携带自己的行末偏移。
pending, end = read_after_cursor() # 读取,不改持久游标
with count_state_lock():
state = load_count_state()
applied = state.applied
for end_offset, key in pending:
if end_offset > applied:
state.counts[key] = state.counts.get(key, 0) + 1
state.applied = max([applied] + [pos for pos, _ in pending])
replace_json_as_one_unit(state) # 计数与水位同次替换
advance_cursor(end) # 完成扫描工作后才推进
这段程序允许同一输入至少一次进入读取路径。它把重复的处理效果限制在计数这一步:end_offset <= A 的普通信号已经计过,恢复时跳过;end_offset > A 的信号才增加累计值。这里的幂等,指同一批旧信号再次输入时,不再改变已经完成的计数。恢复时,每轮重新读取已保存的计数状态;上一轮内存中曾算过多少、执行到哪里,都不参与去重判断。
重读三条,累计仍然是三
现在给两条旧信号标上合成的字节位置。假定每条记录连同行尾恰好占 40 字节,第一条结束于 40,第二条结束于 80。这些数字用于演算,与真实记录长度无关。两条信号都属于 K,初始状态是 N[K]=0,A=0,C=0。
第一轮读出两条记录,把 N[K]=2 和 A=80 一起保存。随后,在写游标前发生故障,C 留在 0。此时新增第三条 K 信号,行末偏移为 120。第二轮从旧游标开始,确实会读出三条记录,读取次数多于新增次数。
图 2 中的水位位于 80。偏移 40 与 80 都满足 end_offset <= A,所以它们进入扫描器,却不再贡献加一;偏移 120 越过水位,才使 N 从 2 增到 3。随后,系统一起保存 N[K]=3,A=120,最后把 C 推到 120。
C=0 从头重读A=80 已计水位end=40 ≤ A跳过计数end=80 ≤ A跳过计数end=120 > A计数 +1N[K]:2 → 3图 3 展开完整的状态轨迹。如果第二轮只是对读出的三条信号再次累加,结果会是 2+3=5;带着水位恢复时,算式是 2+0+0+1=3。两条旧信号被读了两次,但各自的计数效果只有一次。
- 01
初始状态
C=0 · A=0 · N[K]=0 - 02
两条已计,游标推进前失败
两条旧信号行末:40、80。C=0 · A=80 · N[K]=2 - 03
追加一条,恢复时读三条
40+080+0120+1 - 04
保存计数状态,再推进游标
C=120 · A=120 · N[K]=3
行末偏移能承担这个判断,依赖固定文件只追加的前提。这里无需给每条信号分配另一个全局标识,但也继承了文件位置的限制:同样的偏移只有在同一个稳定的文件历史里,才指向同一条记录。
测试在哪个窗口停下
针对这条路径的两个测试,分别让执行停在计数保存前、计数保存后。第一个在首次保存时注入 OSError,断言失败后累计值仍为 0,下一次扫描重新读到两条信号,累计值达到 2。它检查的是“保存失败不能提前消费输入”。
第二个让游标推进函数抛出异常,断言此时计数已保存为 2。恢复游标写入后再追加一条信号,它要求下一次扫描读到三条,累计值却只增长到 3;再扫一次则没有待读信号。它检查的是“结果已保存,但读取进度落后”时,旧输入重放能否保持幂等。
两个测试把同一条恢复规则放到相反的失败窗口里:保存未完成,就保留输入;保存已完成,就凭同次保存的水位跳过旧计数。读取进度与业务效果因此分开:结果与去重证据属于同一个保存单元,游标最后推进。恢复可以重复读取,而累计值仍沿着已完成的效果继续增长。