持久化执行前置课(六):事务发件箱,让状态推进和命令发布不再双写 上一篇让重复活动收敛到同一业务回执。然而决定核一次会同时产生新事件和待执行命令若先更新状态再投队列崩溃会丢命令若先投队列再更新状态消费者可能看到尚不存在的前提。一、痛点两次成功之间总有崩溃窗口数据库事务无法覆盖普通消息代理所谓“按顺序调用两个可靠系统”仍然是不可靠的双写。重试可以修复部分失败却无法判断第一次发布是否已经被代理接受。更麻烦的是消息可能在数据库提交前被极速消费者处理产生一个找不到源事件的孤儿副作用。事务发件箱把命令先作为数据库行与事件在同一事务提交。独立发布器扫描未发布行发送后标记完成。发布器崩溃会导致重复发送而非丢失因此消费者仍需使用上一篇的活动幂等键。这个组合明确给出保证数据库事实与命令意图原子产生传输至少一次业务副作用按键去重。二、原理同库原子跨界重试下面程序在一个 SQLite 事务中追加事件和发件箱记录。故意制造重复逻辑命令时command_id唯一约束拒绝第二行。读取者永远不会看到只有事件而没有命令意图的中间状态。importjsonimportsqlite3 dbsqlite3.connect(:memory:)db.executescript( CREATE TABLE events( run_id TEXT NOT NULL, seq INTEGER NOT NULL, kind TEXT NOT NULL, PRIMARY KEY(run_id, seq) ); CREATE TABLE outbox( command_id TEXT PRIMARY KEY, run_id TEXT NOT NULL, kind TEXT NOT NULL, payload TEXT NOT NULL, published INTEGER NOT NULL DEFAULT 0 ); )defcommit_decision(run_id:str,seq:int,event:str,command_id:str,command:str,payload:dict)-None:withdb:db.execute(INSERT INTO events VALUES(?,?,?),(run_id,seq,event))db.execute(INSERT INTO outbox(command_id,run_id,kind,payload) VALUES(?,?,?,?),(command_id,run_id,command,json.dumps(payload,sort_keysTrue,separators(,,:))),)commit_decision(run-1,1,trip_requested,cmd-1,lock_budget,{amount:900})event_countdb.execute(SELECT COUNT(*) FROM events).fetchone()[0]command_countdb.execute(SELECT COUNT(*) FROM outbox).fetchone()[0]print(events,event_count,commands,command_count)print(db.execute(SELECT command_id,kind,published FROM outbox).fetchone())输出events 1 commands 1 (cmd-1, lock_budget, 0)三、实现发布租约允许崩溃后接管多个发布器需要避免同时长期处理同一行同时又要允许死节点的任务被接管。示例使用逻辑时钟与租约领取时把行改为sending并记录截止点超时行可重新领取确认发布只接受持有该租约的 worker。真实数据库应使用行锁或带版本号的条件更新。fromdataclassesimportdataclassdataclassclassItem:command_id:strstatus:strpendingowner:str|NoneNonelease_until:int0deliveries:int0defclaim(item:Item,worker:str,now:int,ttl:int)-bool:availableitem.statuspendingor(item.statussendinganditem.lease_untilnow)ifnotavailable:returnFalseitem.statussendingitem.ownerworker item.lease_untilnowttl item.deliveries1returnTruedefacknowledge(item:Item,worker:str)-None:ifitem.status!sendingoritem.owner!worker:raiseRuntimeError(lease_not_owned)item.statuspublisheditemItem(cmd-1)print(a_claimed,claim(item,worker-a,now0,ttl5))print(b_early,claim(item,worker-b,now3,ttl5))print(b_takeover,claim(item,worker-b,now6,ttl5))acknowledge(item,worker-b)print(status,item.status,deliveries,item.deliveries)输出a_claimed True b_early False b_takeover True status published deliveries 2四、踩坑已标记发布不代表已处理发件箱的published只证明消息代理接受了消息不能证明活动完成完成必须由独立回执事件表达。若在发送前标记发布会丢消息发送后再标记则必然存在重复窗口这是设计允许的行为。另一个坑是扫描没有索引的整张表积压后会拖垮主库应为状态和可用时间建立索引并限制批量。租约不是锁定业务所有权的永久凭证。worker 发生长暂停后可能在租约过期时继续发送因此下游幂等仍不可省。清理发件箱也要等到消息已发布、完成回执已持久化且保留期满足不能只看到published1就删除。失败消息应进入可检查状态而非无上限热循环。五、验证守恒关系比成功日志可靠对每个会产生命令的事件数据库中必须恰有一个稳定命令 ID任何发件箱行都必须能追溯到源运行完成事件必须引用已存在命令。测试在事务中途抛异常应看到事件和命令都没有提交在发送后、确认前杀死发布器应看到重复投递但只有一个业务结果。到这里事实、命令意图和外部回执已经连成链。下一篇将系统化讨论崩溃恢复怎样从租约、心跳和明确的“不确定”状态中判断该自动接管还是停给人工。参考来源Microsoft事务发件箱模式AWS事务发件箱SQLite原子提交 觉得有用就点个赞 收藏方便回头查阅有疑问直接在评论区留言我看到都会回。 本文属于《持久化执行前置课》系列持续更新关注不迷路。 文章里的代码都能直接跑。想要可直接 clone 的完整工程 配套部署脚本 / 踩坑清单评论一声或发邮件到cj2664qq.com我免费发你。如果你正好在做类似系统、或有工程化难题想找人做也欢迎邮件聊一句——我按实际情况评估能落地的就接单或出方案。评论和邮件都能直接找到我不用跳别的平台。