上一篇用事务发件箱消除了事件与命令意图之间的双写,但 worker 仍可能停在任意指令。恢复器若只看“最后心跳超时”就重跑,会把暂时卡顿、已完成未回报和真正失联混成一种情况。
一、痛点:超时只说明我们不知道
进程消失前可能尚未领取、正在计算、远端已成功、本地已提交。超时不是失败事实,只是观察者在期限内没有得到新证据。可靠系统应把unknown作为一等状态:纯计算可以重新执行,有幂等回执的活动可以查询或重投,不可查询且不可逆的活动则必须暂停人工核对。
二、原理:恢复决策是一张显式表
恢复输入至少包含租约、活动类别、回执和尝试次数。下面程序把规则写成纯函数,避免散落在异常处理里。相同证据永远产生相同动作,manual_review不会被后台循环偷偷改成重试。
fromdataclassesimportdataclass@dataclass(frozen=True)classEvidence:lease_expired:boolreceipt:str|Noneretry_safe:boolattempts:intmax_attempts:intdefrecover(e:Evidence)->str:ifnote.lease_expired:return"wait"ife.receipt=="succeeded":return"record_completion"ife.receipt=="failed_permanent":return"record_failure"ifnote.retry_safe:return"manual_review"ife.attempts>=e.max_attempts:return"exhausted"return"retry_same_activity"cases=[Evidence(False,None,True,1,3),Evidence(True,"succeeded",True,1,3),Evidence(True,None,False,1,3),Evidence(True,None,True,1,3),]foritemincases:print(recover(item))输出:
wait record_completion manual_review retry_same_activity三、实现:租约使用数据库时间和栅栏令牌
仅有租约会遇到旧 worker 复活。每次领取递增fence,提交结果必须携带当前令牌;过期持有者即使恢复,也无法覆盖新持有者。示例用内存对象展示条件更新语义,生产中应放进一个数据库事务。
fromdataclassesimportdataclass@dataclassclassLease:owner:str|None=Noneuntil:int=0fence:int=0result:str|None=Nonedefclaim(lease:Lease,worker:str,now:int,ttl:int)->int:iflease.ownerisnotNoneandlease.until>now:raiseRuntimeError("busy")lease.owner=worker lease.until=now+ttl lease.fence+=1returnlease.fencedefcomplete(lease:Lease,worker:str,fence:int,result:str)->None:iflease.owner!=workerorlease.fence!=fence:raiseRuntimeError("stale_worker")lease.result=result lease=Lease()old=claim(lease,"a",0,5)new=claim(lease,"b",6,5)try:complete(lease,"a",old,"late")exceptRuntimeErroraserror:print(error)complete(lease,"b",new,"ok")print(lease.fence,lease.result)输出:
stale_worker 2 ok四、踩坑:心跳不是进度,接管也不是回滚
worker 可以持续心跳却死循环,因此还要记录最后事件序号和活动阶段。反过来,长时间 API 调用没有心跳也未必失败,租约应按活动类型配置并可续期。不要使用各 worker 本地时钟裁决租约;时钟漂移会制造双主,应依赖数据库时间或协调服务。
恢复器不能删除旧尝试。旧日志、栅栏令牌和回执共同解释为何接管。对于不安全活动,人工确认也应追加结构化事件,写明证据和操作者,而不是直接改状态。这样之后的重放仍能得到相同结论。
五、验证:在每个阶段注入死亡
测试应在领取前后、远端调用前后、回执写入前后和完成提交前后杀死 worker。安全活动最终应完成一次;不可判定活动应稳定停在人工状态;旧栅栏提交必须被拒绝。还要模拟心跳线程存活而业务线程卡死,验证进度监控能报警。
崩溃分类解决“能否接管”,下一篇继续约束“可以重试多少次”。没有预算的安全重试仍会形成流量风暴,并把永久错误伪装成暂时故障。
参考来源
- Martin Kleppmann:分布式锁的栅栏令牌
- Kubernetes:Lease
- Google SRE:处理过载
👍 觉得有用就点个赞 + 收藏,方便回头查阅;有疑问直接在评论区留言,我看到都会回。
🚀 本文属于《持久化执行前置课》系列,持续更新,关注不迷路。
📌 文章里的代码都能直接跑。想要可直接 clone 的完整工程 + 配套部署脚本 / 踩坑清单?评论一声或发邮件到cj2664@qq.com,我免费发你。
如果你正好在做类似系统、或有工程化难题想找人做,也欢迎邮件聊一句——我按实际情况评估,能落地的就接单或出方案。评论和邮件都能直接找到我,不用跳别的平台。