复用 replay.py 的恢复结论
上一篇产出的replay.py会从EventStore.load恢复ReplayResult,并明确让pending_commands永远不从重放结果派发命令。这留下一个必须回答的问题:真正待执行的命令存在哪里?本篇以replay(...).last_seq作为乐观版本,以decide产生的Command作为输入,在 SQLite 中增加事务 outbox。
最危险的窗口是:事件已经提交,但命令还没进入队列,进程恰好崩溃。恢复时状态认为命令已产生,外部世界却永远收不到它。反过来,先发命令再提交事件,外部成功而本地崩溃,重试就可能执行两次。单个数据库无法与任意 HTTP 服务做原子事务,因此工程上采用“本地原子、外部幂等”:事件和 outbox 在一个 SQLite 事务提交;worker 至少一次投递;接收方用稳定幂等键把重复执行折叠为一次。
把事件与命令写在一个提交里
保存为durable_store.py。命令 ID 由工作流、触发事件序号和命令下标确定性构造,不能每次重试生成 UUID,否则同一逻辑命令会得到不同身份,接收方无法去重。processed_at为空表示待投递,成功后再标记。真实系统还应记录尝试次数和错误,本篇先聚焦原子性。
importjsonimportsqlite3fromevent_storeimportEventStorefromminiflowimportCommand OUTBOX_SCHEMA=""" CREATE TABLE IF NOT EXISTS outbox ( command_id TEXT PRIMARY KEY, workflow_id TEXT NOT NULL, event_seq INTEGER NOT NULL, kind TEXT NOT NULL, payload TEXT NOT NULL, processed_at TEXT ); """classDurableStore(EventStore):def__init__(self,path:str):super().__init__(path)withself.connect()asdb:db.executescript(OUTBOX_SCHEMA)defappend_with_commands(self,workflow_id:str,expected_seq:int,kind:str,payload:dict,commands:list[Command])->int:seq=expected_seq+1withself.connect()asdb:db.execute("INSERT INTO events(workflow_id,seq,kind,payload) VALUES(?,?,?,?)",(workflow_id,seq,kind,json.dumps(payload,ensure_ascii=False,sort_keys=True)),)forindex,commandinenumerate(commands):command_id=f"{workflow_id}:{seq}:{index}"db.execute("INSERT INTO outbox(command_id,workflow_id,event_seq,""kind,payload) VALUES(?,?,?,?,?)",(command_id,workflow_id,seq,command.kind,json.dumps(command.data,sort_keys=True)),)returnseqdefpending(self)->list[sqlite3.Row]:withself.connect()asdb:returndb.execute("SELECT * FROM outbox WHERE processed_at IS NULL ""ORDER BY workflow_id,event_seq,command_id").fetchall()defmark_done(self,command_id:str)->None:withself.connect()asdb:db.execute("UPDATE outbox SET processed_at=CURRENT_TIMESTAMP ""WHERE command_id=?",(command_id,))运行输出:
(模块定义成功,无标准输出)这里必须强调,mark_done和远端 HTTP 成功之间仍存在窗口:远端成功后、本地标记前崩溃,命令会重投。因此 outbox 提供的是至少一次,不是恰好一次。所谓“端到端恰好一次”通常需要接收方参与:它把command_id放入唯一约束,并在自己的业务事务中同时写去重记录和业务结果。
用本地接收方证明重复被折叠
保存为demo_104.py。FakeSupplier模拟支持幂等键的供应商。我们故意在第一次远端成功后不调用mark_done,模拟确认前崩溃;第二次扫描 outbox 会再次投递相同命令 ID,但供应商返回原结果,不产生第二张订单。
importtempfilefrompathlibimportPathfromdurable_storeimportDurableStorefromminiflowimportCommandclassFakeSupplier:def__init__(self):self.results={}self.side_effect_count=0defexecute(self,command_id:str,kind:str,payload:str)->str:ifcommand_idinself.results:returnself.results[command_id]self.side_effect_count+=1result=f"REF-{self.side_effect_count}"self.results[command_id]=resultreturnresultwithtempfile.TemporaryDirectory()asdirectory:store=DurableStore(str(Path(directory)/"flow.db"))supplier=FakeSupplier()store.append_with_commands("trip-001",0,"trip_requested",{"city":"成都"},[Command("lock_budget",{"trip_id":"trip-001"})],)first=store.pending()[0]print("first delivery:",supplier.execute(first["command_id"],first["kind"],first["payload"]))# 模拟此处崩溃:没有 mark_done。again=store.pending()[0]print("retry delivery:",supplier.execute(again["command_id"],again["kind"],again["payload"]))store.mark_done(again["command_id"])print("real side effects:",supplier.side_effect_count)print("pending:",len(store.pending()))assertsupplier.side_effect_count==1运行输出:
first delivery: REF-1 retry delivery: REF-1 real side effects: 1 pending: 0幂等不等于请求内容相同
常见错误是用请求体哈希当幂等键。两个不同旅客可能提交完全相同的航班请求,它们是两个合法业务操作,不该合并;同一操作重试时,时间戳或追踪字段又可能变化,哈希反而不同。幂等键表达“业务意图的身份”,本例由工作流位置确定。接收方还应存储该键首次请求的参数摘要;如果同一个键后来携带不同核心参数,应返回冲突而不是复用旧结果。这能暴露调用方错误,避免静默关联错误订单。
另一个非平凡坑是把 outbox 行删掉。删除会失去审计依据,也让延迟到达的重复响应难以解释。更稳妥的是标记完成并设置保留期,定期归档;去重记录的保留期必须覆盖生产者可能重试的最长时间。如果接收方 24 小时删除幂等键,而生产者能在七天后重试,“一次”就会重新变成“两次”。
事务 outbox 也不是无限队列。某个永久失败命令会不断重试,需要指数退避、最大尝试策略和人工处理入口;但不能简单标记成功。第六篇会用持久化定时器实现退避,第七篇会区分可重试失败和触发补偿的终局失败。在此之前,最重要的记忆点是确认语义:只有拿到可验证的远端结果,才能mark_done;网络超时意味着结果未知,应该带同一幂等键查询或重试。
本篇交接
本篇产出的durable_store.py新增append_with_commands、pending、mark_done,并确定了trip-001:事件序号:命令下标形式的稳定command_id。下一篇会直接复用 outbox 行和命令 ID,但不允许所有 worker 同时执行它们:我们将增加带过期时间的租约,让多个进程安全竞争,同时处理持有者暂停、崩溃和迟到确认。
👍 觉得有用就点个赞 + 收藏,方便回头查阅;有疑问直接在评论区留言,我看到都会回。
📌 文章里的代码都能直接跑。想要可直接 clone 的完整工程 + 配套部署脚本 / 踩坑清单?评论一声或发邮件到cj2664@qq.com,我免费发你。
如果你正好在做类似系统、或有工程化难题想找人做,也欢迎邮件聊一句——我按实际情况评估,能落地的就接单或出方案。评论和邮件都能直接找到我,不用跳别的平台。