执行成功,确认却丢失了

任务要求把累计值增加 7。处理器完成加法,确认消息随后丢失;发送方无法确定结果,于是再次执行任务。没有幂等保护时,两次尝试把累计值变成 14,即使第二次正常返回,也不能据此认定业务只执行过一次。

本篇用 Ruby 3.4.11 的 Queue、Mutex 和 Thread 建立本地替身。它真实记录执行次数、累计值和事件顺序,故障由固定条件注入,不依赖网络偶发失败。完整源码位于 examples/ruby/labs/E08/,附件为 本篇实验源码。

前置内容是异常、测试和并发。实验不连接 Redis,也不安装任务框架;队列、幂等记录和失败事件都仅存在于当前进程内。目标是先明确协议需要保证什么,再讨论框架能力是否覆盖这些要求。

返回值、确认和副作用是三个事件

一次处理可以分成收到任务、执行副作用、确认完成三段。失败发生在不同位置,对重试的要求不同:副作用前失败通常尚未改变目标;副作用后但确认前失败,则结果可能已经存在。

sequenceDiagram
  participant W as Worker
  participant S as EffectStore
  W->>S: apply_once(a, 7)
  S-->>W: applied,总值增加 7
  Note over W: 确认丢失,记录 retry
  W->>S: apply_once(a, 7)
  S-->>W: duplicate,总值不变
  Note over W: acked

lab 先运行没有幂等记录的反例,断言两次尝试得到 14。随后注入同样的确认丢失,检查受保护实现只写入一次。对外部消息系统,重复投递是协议设计时需要考虑的情况;例如 SQS 标准队列明确要求消费者能够容忍重复。SQS 至少一次投递

这项实验没有模拟 SQS,也不替任何框架验证端到端交付保证。它只把“执行完成但上游仍然重试”这一状态暴露为可复现输入。验收对象是业务累计值,不能只看 worker 最终返回了什么。

幂等键必须绑定同一个业务操作

Job 保存 key、amount、故障次数和是否丢失确认。提交时复制并冻结 key,避免调用者提交后再改字符串导致去重身份漂移。key 有固定 ASCII 格式与长度,amount 必须是正整数。

EffectStore 用同一个 Mutex 保护检查、写入和记录:

1
2
3
4
5
6
7
8
9
10
@mutex.synchronize do
if @seen.key?(key)
raise KeyConflict unless @seen[key] == amount
return :duplicate
end
@total += amount
@writes += 1
@seen[key] = amount
:applied
end

同键同参数表示同一操作,可以返回已经处理;同键不同参数则明确报错。静默忽略不同 payload 会掩盖调用方复用键的错误,重新执行又会破坏幂等语义。业务系统通常还应在键中纳入操作种类和业务标识,并定义记录保留期。

互斥区不能只包住 @seen 查询。两个线程若都查到不存在,再分别执行副作用,仍会重复。本实验用 Queue 同时释放两个调用线程,最后断言一个 applied、一个 duplicate,实际写入次数为 1。Mutex 保护的是共享 Ruby 状态,语义来自受保护区域,而不是来自碰巧存在的 GVL。Ruby Mutex

此处累计值和幂等记录都在同一个内存锁中,原子范围清楚。若把加法换成 HTTP 请求,远端副作用不会因为本地 Mutex 而与本地记录原子提交。进程也可能在远端成功后、记录写入前崩溃。届时应由接收端接受幂等键,或采用适合实际资源的持久协议;不能将此类直接包装成生产消息框架。

只重试定义过的失败,并给出终点

Worker 默认最多尝试三次。fail_before 模拟副作用前的临时失败,lose_ack 模拟副作用后的确认丢失。两者都使用显式的 Retryable,处理器只捕获这类错误。键冲突和未知程序错误不会被当作临时故障无限重试。

重试间隔采用 2**(attempt - 1),三次尝试之间对应 1、2 两个等待值。实验通过注入 sleeper 记录计划时间,不实际睡眠,因此验证的是调度参数与次数,不是实际延迟或吞吐量。真实实现还要决定最大等待、抖动、总截止时间和哪些错误值得重试。

队列依次接收:丢失首次确认的 a、重复投递的 a、持续失败的 poison、正常的 b。a 金额 7,b 金额 3,poison 金额 99。最终输出应为:

1
2
3
naive_total=14
deduplicated={total: 10, writes: 2, keys: ["a", "b"]}
delays=[1, 1, 2]

第一个等待属于 a 的确认丢失,后两个属于 poison。poison 恰好尝试三次后留下 dead 事件,未产生副作用;b 随后仍能完成。重试在当前 worker 内串行进行,所以毒任务会占用后续任务的等待时间。这是教学实现的调度选择,生产系统可以使用延迟队列,但必须重新验证顺序和关闭规则。

这里的 dead 只是内存中的终止记录,没有持久死信队列、检索页面或人工重放流程。把失败记录打印出来,不等于已经具备长期恢复能力。恢复时还必须保留原始幂等身份,否则“重放”可能变成另一项新业务操作。

取消必须在明确的检查点发生

取消 token 通过 Mutex 共享一个布尔状态。Worker 在每次尝试开始时检查一次,在测试 hook 返回后、执行副作用前再检查一次。取消测试用两个 Queue 建立屏障:

1
2
3
4
entered.pop
cancel_worker.token.cancel
release << true
cancel_worker.close(drain: false)

worker 到达 before_effect 后先通知 entered,再等待 release。主线程确认它已经到达检查点,设置取消,最后释放。由此可以稳定断言当前任务和待处理任务都记录 cancelled,累计写入仍为零。使用固定睡眠代替屏障,只能提高某种时序发生的概率,不能证明取消在副作用之前。

取消是协作协议,无法在任意机器指令之间保证“零副作用”。即使已经完成一次检查,取消也可能随后到达。本实现没有把检查 token 与外部副作用做成一个原子操作,因此不提供抢占式取消保证。

lab 还单独演示晚到取消:先写入金额 4,再设置 token,最后累计值仍为 4。要撤销已经发生的操作,需要业务补偿协议;改变线程状态不会自动生成补偿。不能把 cancelled? 为真当成目标状态已经恢复。

有序关闭先禁止提交,再处理剩余工作

close 在状态锁内禁止新任务提交,然后关闭 Queue。Ruby 的关闭队列仍允许取出已有元素;当队列已经空,阻塞形式的 pop 返回 nil。因此循环可以在排空后自然退出。Queue#close

本实现拒绝 nil 任务,使 nil 只表示队列已结束。若业务允许把 nil 入队,原本的退出条件就会误把正常内容当成关闭信号。结束标志应与合法任务空间分开定义。

默认 drain: true 保留既有任务的处理流程,包括有上限的重试。事件里 b 的 acked 必须先于 closed。drain: false 则设置共享取消 token,尚未越过检查点的任务跳过副作用;两种方式最终都等待线程结束。

实现使用 Thread#value 等待 worker,也让未捕获异常回到调用者。只调用 close 后立刻让主进程退出,无法证明处理线程已经完成清理。Ruby Thread

本实验的 hook 与 sleeper 必须协作返回。若一个 hook 永久阻塞,关闭也会等待;代码没有强杀线程、关闭期限或跨进程监督。一个本地替身不应通过省略这些事实让读者误以为它已经解决了任意运行任务的停止问题。

队列还存在容量边界。本篇使用普通 Queue,没有生产者背压;若生产速度长期高于处理速度,内存中的待处理任务就会增长。迁移到 SizedQueue 或外部消息系统时,应重新设计提交超时与关闭行为,而不是把“队列操作线程安全”解释为整个系统容量安全。

幂等记录同样会增长。教学用例只有少量固定键,未实现过期或清理。生产系统若让记录过早失效,晚到的重复任务可能重新执行;若永久保存,又需要计算存储成本。保留窗口必须与业务重试、人工重放和消息保留期限一起定义,不能只选择一个看起来方便的缓存时间。

事件读取通过锁取得快照,最终断言放在 worker 结束以后,避免把运行中的中间状态当作完成状态。在线监控当然需要读取中间状态,但应把“正在重试”“已停止”“业务成功”分开表示。本例的 closed 仅证明循环已经结束,是否成功仍需检查具体任务和副作用记录。

用事件与业务状态共同验收

运行命令为:

1
ruby examples/ruby/labs/E08/run.rb

正常日志最后是 PASS lab E08。断言覆盖未保护的重复加法、两次受保护写入、重复投递、键冲突、三次重试上限、后续任务排空、关闭后拒绝提交、并发同键与取消场景。

事件用于解释路径,累计值和写入次数用于验证效果。只检查 acked 会漏掉重复写入,只检查总值又可能把一笔漏写与另一笔多写抵消。因此实验还比较写入数、已处理键集合和每个关键事件。

协议问题 本地观察 未被证明的内容
确认丢失后重投 a 的重复返回不增加总值 网络消息服务的交付保证
两线程同键 一次 applied、一次 duplicate 跨进程与崩溃后的原子性
达到上限 poison 三次后 dead 持久死信和人工恢复
取消和关闭 检查点前零写入,线程已结束 任意时刻强制终止

练习一:把幂等检查与写入拆到两个互斥区之间,用 Queue 屏障让两个线程都读到“不存在”,写出能够稳定观察两次副作用的失败实验,再恢复同一临界区。

练习二:为任务增加最大累计等待预算。让重试上限仍是三次,但等待预算只允许一次延迟,断言任务提前进入明确的终止状态,后续任务仍可排空。不要把所有异常都改成 Retryable 来简化实现。

参考资料

系列导航

导读 · 上一篇:E07:受限规则语言:从 block DSL 到独立解析器 · 完整源码包