持久化执行前置课(六):事务发件箱,让状态推进和命令发布不再双写
2026/8/2 5:16:41 网站建设 项目流程

上一篇让重复活动收敛到同一业务回执。然而决定核一次会同时产生新事件和待执行命令:若先更新状态再投队列,崩溃会丢命令;若先投队列再更新状态,消费者可能看到尚不存在的前提。

一、痛点:两次成功之间总有崩溃窗口

数据库事务无法覆盖普通消息代理,所谓“按顺序调用两个可靠系统”仍然是不可靠的双写。重试可以修复部分失败,却无法判断第一次发布是否已经被代理接受。更麻烦的是消息可能在数据库提交前被极速消费者处理,产生一个找不到源事件的孤儿副作用。

事务发件箱把命令先作为数据库行,与事件在同一事务提交。独立发布器扫描未发布行,发送后标记完成。发布器崩溃会导致重复发送而非丢失,因此消费者仍需使用上一篇的活动幂等键。这个组合明确给出保证:数据库事实与命令意图原子产生,传输至少一次,业务副作用按键去重。

二、原理:同库原子,跨界重试

下面程序在一个 SQLite 事务中追加事件和发件箱记录。故意制造重复逻辑命令时,command_id唯一约束拒绝第二行。读取者永远不会看到只有事件而没有命令意图的中间状态。

importjsonimportsqlite3 db=sqlite3.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_keys=True,separators=(",",":"))),)commit_decision("run-1",1,"trip_requested","cmd-1","lock_budget",{"amount":900})event_count=db.execute("SELECT COUNT(*) FROM events").fetchone()[0]command_count=db.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。真实数据库应使用行锁或带版本号的条件更新。

fromdataclassesimportdataclass@dataclassclassItem:command_id:strstatus:str="pending"owner:str|None=Nonelease_until:int=0deliveries:int=0defclaim(item:Item,worker:str,now:int,ttl:int)->bool:available=item.status=="pending"or(item.status=="sending"anditem.lease_until<=now)ifnotavailable:returnFalseitem.status="sending"item.owner=worker item.lease_until=now+ttl item.deliveries+=1returnTruedefacknowledge(item:Item,worker:str)->None:ifitem.status!="sending"oritem.owner!=worker:raiseRuntimeError("lease_not_owned")item.status="published"item=Item("cmd-1")print("a_claimed=",claim(item,"worker-a",now=0,ttl=5))print("b_early=",claim(item,"worker-b",now=3,ttl=5))print("b_takeover=",claim(item,"worker-b",now=6,ttl=5))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 发生长暂停后可能在租约过期时继续发送,因此下游幂等仍不可省。清理发件箱也要等到消息已发布、完成回执已持久化且保留期满足,不能只看到published=1就删除。失败消息应进入可检查状态,而非无上限热循环。

五、验证:守恒关系比成功日志可靠

对每个会产生命令的事件,数据库中必须恰有一个稳定命令 ID;任何发件箱行都必须能追溯到源运行;完成事件必须引用已存在命令。测试在事务中途抛异常,应看到事件和命令都没有提交;在发送后、确认前杀死发布器,应看到重复投递但只有一个业务结果。

到这里,事实、命令意图和外部回执已经连成链。下一篇将系统化讨论崩溃恢复:怎样从租约、心跳和明确的“不确定”状态中判断该自动接管还是停给人工。

参考来源

  • Microsoft:事务发件箱模式
  • AWS:事务发件箱
  • SQLite:原子提交

👍 觉得有用就点个赞 + 收藏,方便回头查阅;有疑问直接在评论区留言,我看到都会回。

🚀 本文属于《持久化执行前置课》系列,持续更新,关注不迷路。

📌 文章里的代码都能直接跑。想要可直接 clone 的完整工程 + 配套部署脚本 / 踩坑清单?评论一声或发邮件到cj2664@qq.com,我免费发你。
如果你正好在做类似系统、或有工程化难题想找人做,也欢迎邮件聊一句——我按实际情况评估,能落地的就接单或出方案。评论和邮件都能直接找到我,不用跳别的平台。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询