事务 Outbox 与幂等键:跨过“提交后崩溃”的缝隙 复用 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:seqexpected_seq1withself.connect()asdb:db.execute(INSERT INTO events(workflow_id,seq,kind,payload) VALUES(?,?,?,?),(workflow_id,seq,kind,json.dumps(payload,ensure_asciiFalse,sort_keysTrue)),)forindex,commandinenumerate(commands):command_idf{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_keysTrue)),)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_atCURRENT_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_count0defexecute(self,command_id:str,kind:str,payload:str)-str:ifcommand_idinself.results:returnself.results[command_id]self.side_effect_count1resultfREF-{self.side_effect_count}self.results[command_id]resultreturnresultwithtempfile.TemporaryDirectory()asdirectory:storeDurableStore(str(Path(directory)/flow.db))supplierFakeSupplier()store.append_with_commands(trip-001,0,trip_requested,{city:成都},[Command(lock_budget,{trip_id:trip-001})],)firststore.pending()[0]print(first delivery:,supplier.execute(first[command_id],first[kind],first[payload]))# 模拟此处崩溃没有 mark_done。againstore.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_count1运行输出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 的完整工程 配套部署脚本 / 踩坑清单评论一声或发邮件到cj2664qq.com我免费发你。如果你正好在做类似系统、或有工程化难题想找人做也欢迎邮件聊一句——我按实际情况评估能落地的就接单或出方案。评论和邮件都能直接找到我不用跳别的平台。