Harness 工程进展:合并、事件与恢复 · 2026年10月1日 · 10 分钟

游标最后推进:从两次故障推导信号重放

一批信号可以重读,却不能重复加到累计值里。沿着丢信号与重复计数两个故障窗口,推导游标、计数和字节水位的写入顺序。

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 把这两个保存边界放在同一条时间轴上。按这个顺序,进程在两次保存之间退出,会留下“计数状态已更新,读取游标未更新”的组合,水位足以解释它;“读取游标已更新,计数状态未更新”则不应由正常路径产生。

  1. 01

    读入内存,保留旧游标

    C=0 · N[K]=0 · A=0

  2. 02

    先保存计数与水位

    N[K]=2 · A=80

    flock 包住完整读改写。
  3. ×

    进程在此退出

    C=0 · N[K]=2 · A=80

    下轮重读,旧信号不再加一。
  4. 03

    最后保存读取游标

    C=80 · N[K]=2 · A=80

C 是读取游标,A 是已计水位;计数状态与游标分两次保存。
图 1|进程在两次保存之间退出,会留下旧游标与新计数状态。下一轮会重读,但已计水位让这些输入不再增加累计值。

下面的伪代码只展示有效普通信号的计数分支,保留决定恢复行为的顺序。信号列表在读取阶段已按文件顺序得到,每条携带自己的行末偏移。

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 已计水位
旧信号 1 · Kend=40 ≤ A跳过计数
旧信号 2 · Kend=80 ≤ A跳过计数
新信号 3 · Kend=120 > A计数 +1
重读 3 条,新增计数 1 次N[K]:2 → 3
仅展示有效的普通计数信号;40 / 80 / 120 为合成字节偏移。
图 2|读取从 C=0 开始,计数从 A=80 之后继续。游标落后不等于计数遗漏,二者之间的旧记录可以安全重放。

图 3 展开完整的状态轨迹。如果第二轮只是对读出的三条信号再次累加,结果会是 2+3=5;带着水位恢复时,算式是 2+0+0+1=3。两条旧信号被读了两次,但各自的计数效果只有一次。

  1. 01

    初始状态

    C=0 · A=0 · N[K]=0

  2. 02

    两条已计,游标推进前失败

    C=0 · A=80 · N[K]=2

    两条旧信号行末:40、80。
  3. 03

    追加一条,恢复时读三条

    40+080+0120+1
  4. 04

    保存计数状态,再推进游标

    C=120 · A=120 · N[K]=3

2 + 0 + 0 + 1 = 3直接重复累加会得到 2 + 3 = 5。
图 3|故障前已经保存的 2,与恢复时新计入的 1,共同形成最终的 3;游标在最后一次写入后追上水位。

行末偏移能承担这个判断,依赖固定文件只追加的前提。这里无需给每条信号分配另一个全局标识,但也继承了文件位置的限制:同样的偏移只有在同一个稳定的文件历史里,才指向同一条记录。

测试在哪个窗口停下

针对这条路径的两个测试,分别让执行停在计数保存前、计数保存后。第一个在首次保存时注入 OSError,断言失败后累计值仍为 0,下一次扫描重新读到两条信号,累计值达到 2。它检查的是“保存失败不能提前消费输入”。

第二个让游标推进函数抛出异常,断言此时计数已保存为 2。恢复游标写入后再追加一条信号,它要求下一次扫描读到三条,累计值却只增长到 3;再扫一次则没有待读信号。它检查的是“结果已保存,但读取进度落后”时,旧输入重放能否保持幂等。

两个测试把同一条恢复规则放到相反的失败窗口里:保存未完成,就保留输入;保存已完成,就凭同次保存的水位跳过旧计数。读取进度与业务效果因此分开:结果与去重证据属于同一个保存单元,游标最后推进。恢复可以重复读取,而累计值仍沿着已完成的效果继续增长。