双写问题发生在两个提交点之间
业务事务与消息代理各自只能保证自己的提交。若先写订单再发送事件,数据库成功而代理暂时不可用会让下游永远不知道订单存在;若先发送再提交订单,数据库回滚后下游会处理一条幽灵事件。把发送调用放进数据库事务也没有消除问题:远端确认可能在本地回滚前返回,网络超时还会让结果未知,同时长事务持续占用连接和行锁。Outbox 的目标不是神奇的跨系统原子,而是把必须同时成立的本地事实放进同一事务。
本地事务写业务聚合和一条待投递 outbox 记录,两者要么一起提交,要么一起回滚。独立发布器随后读取 outbox 并投递,成功确认后标记结果。这个流程通常提供至少一次投递,因此消费者幂等是协议的一部分,不是可选优化。设计评审要明确事件在哪个业务状态生成、谁拥有 schema、允许怎样重复、排序范围是什么,以及数据库已提交但消息尚未出现时下游能容忍多长延迟。
事件记录包含稳定身份和业务版本
outbox 行至少保存 event_id、aggregate_type、aggregate_id、aggregate_version、event_type、payload、created_at 与投递状态。event_id 在首次业务事务中生成,任何重试都不能换新值;消费者用它识别重复。aggregate_version 表达同一聚合的顺序,不能拿数据库自增 id 推断跨聚合业务先后。payload 是对外契约,只放消费者需要的字段,避免把数据库整行或敏感用户资料复制到更多系统。
事件 schema 应携带明确版本,并采用兼容演进。新增可选字段通常比重命名或改变语义安全;需要破坏性变化时使用新事件类型或版本,让旧消费者有迁移窗口。payload 构造必须使用事务内即将提交的值,而不是提交后重新查询当前行,否则快速发生的下一次修改可能让旧事件携带新快照。审计数据和集成事件用途不同,不能把不可变审计日志直接当成下游稳定 API。
业务写入与 outbox 插入共享一个数据库事务
服务处理命令时先验证身份、输入和当前聚合版本,在 Prisma 或 SQL 事务里执行条件更新,然后插入对应 outbox。任何一步影响行数不符或违反约束都回滚。不要在事务提交后由应用内存补插事件,也不要依赖 ORM 生命周期钩子在不明确的事务上下文运行。若一个命令产生多个事件,它们的版本和顺序要在事务中确定,并通过唯一约束阻止相同聚合版本重复创建。
事务要保持短小,payload 序列化可以使用已经校验的内存对象,但不能在锁内等待网络。大型正文通常不适合复制进 outbox,可保存不可变资源版本引用;消费者读取引用时要确保目标仍存在且权限边界明确。若下游必须得到当时快照,就保存必要字段并限制大小。事务失败日志记录命令和聚合标识,不要预先宣布“事件已发送”,因为此时 outbox 也没有提交。
发布器的领取、发送和确认分成可恢复阶段
多个发布器可以用短事务领取待发送行,写入租约或使用 SKIP LOCKED,提交后才调用代理。发送成功返回可验证消息标识后,再以 event_id、版本和领取者条件更新为已投递。发送时数据库事务不能一直持锁,否则代理变慢会拖垮业务连接池。若进程在发送后确认前崩溃,租约到期会再次投递同一 event_id,这正是消费者必须幂等的原因。
确认失败属于结果未知,不能把事件删除。记录尝试次数、最后错误类别和下一次时间,暂时故障使用有上限退避。确定性的 schema 拒绝或超大消息进入暂停队列,继续重试不会改善。是否允许后续同聚合事件越过失败事件取决于业务排序:账户余额等严格序列通常必须阻塞并告警,独立通知可能允许继续。策略要显式编码,不能由数据库偶然扫描顺序决定。
消费者用收件箱或业务约束实现幂等
消费者接到消息后,在自己的数据库事务中尝试插入 processed_event(event_id),只有首次插入者执行本地业务变更,并与处理标记一起提交。若业务本身有稳定唯一键,也可以由唯一约束裁决,但仍要能区分重复和真正冲突。先写处理标记、事务外再改业务会在中途失败时丢事件;先改业务、最后写标记则可能重复,二者必须属于同一原子单元。
涉及外部副作用的消费者还要继续传递 event_id 派生的操作键。例如邮件服务若不支持幂等,消费者需要本地发送账本与结果未知状态,不能因为消息已确认就假定邮件只发一次。处理成功后才向代理确认;重试相同消息时返回已有业务结果。processed_event 的保留期至少覆盖代理最大重投和灾难恢复窗口,过早清理会让迟到副本重新执行。
积压和失败要按事件年龄观测
最关键的运行指标是最老未投递事件年龄、待发送数量、每类错误、投递到消费延迟和重复命中率。仅统计每秒发送量会在生产者停摆时显示漂亮的零延迟。指标按 event_type 和依赖分组,避免一个持续失败的大消息掩盖其他流量。日志关联 event_id、aggregate_id、聚合版本、尝试次数与代理响应标识,不记录完整 payload 或认证头。
告警阈值来自业务允许的不一致窗口。若搜索索引可延迟十分钟,通知只能延迟一分钟,两者不应共用一条队列年龄告警。发布器扩容前检查数据库候选查询、代理限流和消费者能力,单纯增加并发可能扩大 429 与重试风暴。暂停事件需要后台修复入口,修复 schema 或内容后以同一 event_id 恢复,不能复制一条新事件让原失败证据失联。
故障注入验证每个崩溃窗口
集成测试使用真实数据库创建业务命令,断言业务行与 outbox 始终同时出现或同时缺失。发布器测试分别在领取提交前、领取后发送前、代理已接受后、本地确认前和确认后强制退出,再启动接管者。最终允许发送调用重复,但消费者业务结果只能存在一份,event_id 和聚合版本保持不变。只 mock 一个始终成功的消息客户端无法验证 Outbox 的核心价值。
并发测试让多个发布器领取同一批事件,断言没有事件永久遗漏;让某条内容无效,确认它暂停但策略允许的其他事件继续。消费者测试在本地业务更新后制造事务失败,processed_event 也必须回滚。恢复测试从历史备份还原生产者数据库,再接入包含较新消息的代理副本,确认重复事件被安全吸收,避免灾难恢复把旧 outbox 当新动作重放。
清理与迁移不能破坏可追踪性
已投递 outbox 会持续增长,清理应依据确认时间和审计保留要求按小批次执行。删除前确认代理重投窗口、消费者对账和灾难恢复均不再需要,必要信息可汇总到长期审计而非保留巨大 payload。清理任务本身用稳定游标和上限,避免一个长 DELETE 产生表膨胀与复制延迟。分区可以帮助生命周期管理,但分区键必须配合候选查询和唯一约束设计。
从“事务内直接发送”迁移时,先部署能写 outbox 的生产者并让发布器影子读取,对比预期事件但不发送;然后灰度真实投递,最后移除旧发送路径。切换窗口要防止两个路径同时对同一命令发送,可用功能开关绑定唯一发布协议。回滚应用时仍要保留新表和兼容 schema,确保已经提交的事件继续可处理。最终运行手册应包含积压恢复、暂停修复、代理不可用和密钥轮换步骤。
分区、顺序键与重放工具必须保持事件身份
消息代理的分区键决定同一聚合事件是否保持发送顺序。以 aggregate_id 作为键可以让单个订单版本有序,但不同订单没有全局顺序,消费者不能用到达时间推断跨聚合因果。需要跨聚合事务视图时应设计专门协调模型,不让 Outbox 承担它无法保证的全局序列。发布器记录代理分区与 offset,便于排查而不把它们当业务版本。
消费者发现版本跳跃时可暂存短时间等待缺失事件,超过窗口后告警并查询事实来源;直接丢弃较新事件会造成永久停滞,按到达顺序覆盖又可能倒退状态。事件处理函数根据 aggregateVersion 明确接受重复、下一个版本和不可恢复间隙。快照同步也携带版本,修复后从确定位置继续,而不是清空整条消费记录。
运维重放工具只选择精确 event_id 或受审计时间范围,默认 dry-run 显示事件类型、数量、消费者和可能副作用。重放保持原 event_id,使幂等消费者吸收已处理消息;若业务明确需要重新执行,应创建具有新意图和新审计的补偿命令,不能偷偷换 event_id 绕过收件箱。高风险重放需要双人确认与速率上限。
Outbox 表与业务表共同备份,恢复旧快照后可能重新出现已经被代理发送的记录。发布器再次投递原 event_id 是预期行为,消费者的保留窗口必须覆盖这种灾难恢复。若清理 processed_event 早于可恢复备份年龄,旧事件会被当成首次处理;备份保留、outbox 清理和消费者收件箱 TTL 因而必须统一规划。
实现片段
await tx.outbox.create({ data: { eventId, aggregateId, payload } })
评论 · 0