Outbox、CDC 与领域事件:数据库事实如何可靠进入事件流
把领域事件生成、Outbox 同事务写入、CDC 发布、去重、清理、契约版本和消费者演进连成完整闭环。
Debezium 官方 Outbox Event Router 以 outbox 行插入为事件来源,并建议用 aggregate id 作为 Kafka key 维持同聚合分区顺序。
双写失败来自两个资源没有共同提交点
先提交数据库再发送消息,进程可能在两者之间崩溃,业务事实永久缺少事件;先发消息再提交数据库,消费者可能看到最终回滚的事实;捕获异常并重试仍无法覆盖断电。Outbox 把业务变更和一条待发布事件记录写入同一个数据库事务,使“业务成立但没有可发布证据”的窗口消失。
Outbox 行通常包含 eventId、aggregateType、aggregateId、eventType、payload、schemaVersion、occurredAt 与 headers。eventId 必须在业务事务内确定,发布器重试不能生成新 id。aggregateId 可作为分区 key,使同一聚合的事件保持局部顺序。payload 保存已经发生的事实,而不是要求发布器再次读取当前业务表;否则发布延迟期间状态变化会让历史事件内容漂移。
Relay 与 CDC 都是至少一次发布器
轮询 relay 通过索引扫描未发布行、锁定小批次、发送并标记完成;要处理多实例抢占、长事务、空轮询和清理。若发送成功后标记前崩溃,消息会重复。CDC 读取数据库变更日志,减少业务表轮询并保留提交顺序,但仍存在连接器重启、位点恢复、Topic 写入与 schema 变更等运维边界。
Debezium outbox router 能把列映射到事件 key、header 和 payload,却不会替业务定义事件语义。Outbox 表的 UPDATE 可能不符合只捕获 INSERT 的路由假设;清理任务不能早于 CDC 位点与审计保留要求。发布延迟、最老未发布行年龄、重复率和清理积压必须可观测。
领域事件不应泄漏内部对象图
领域对象可以包含惰性集合、循环引用和内部状态,直接序列化会把实现细节变成公共契约。集成事件应由领域事实投影而来,只包含消费者需要的稳定字段。字段新增优先可选并提供默认值;枚举新增要求消费者按 unknown 分支处理;重命名和类型收窄通常是破坏性变更,应通过新字段双写、双读和弃用窗口迁移。
Schema Registry 可以检查兼容规则,但语法兼容不等于语义兼容。把 amount 从“分”改成“元”即使仍是 long,也会悄悄破坏消费者。契约评审要记录单位、时区、空值含义、唯一性、顺序键与删除语义,并用旧消费者读取新事件、新消费者读取历史事件的双向测试验证。
用所有权转移图读懂可靠性
失败矩阵比“保证不丢”更可执行
发布前崩溃由本地事务或 Outbox 恢复;写出后响应前断线属于结果未知,需要稳定 id 重发;Broker 接收后副本失效取决于持久化与复制策略;消费后提交前崩溃可以安全重投;业务提交后确认前崩溃会重复;确认后才发现异步副作用失败则进入补偿。每个窗口都要有唯一责任方、持久证据、重试上限和人工退出路径。
工程选择的核心是把不可避免的不确定窗口转成可重复、可查询、可补偿的状态。不要用无限重试掩盖未知结果,也不要用提前确认换取表面低 lag。
Outbox 自己也需要容量与清理协议
高峰期每个业务事务增加一行 Outbox,会扩大 WAL、索引和主从复制压力。表应按提交时间或业务分片建立可持续扫描索引,Relay 小批领取并设置锁等待上限,不能用无界全表扫描追赶积压。发布吞吐低于写入吞吐时,最老未发布年龄比总行数更能说明业务延迟。
已发布行不能在发送回调后立刻删除:回调可能早于 CDC 消费位点或审计落盘。清理水位应综合发布状态、连接器位点、保留期和备份恢复目标,并按小批删除避免长事务。数据库从备份恢复后还可能重新出现已发布 Outbox 行,因此 eventId 去重必须跨恢复周期成立。
用两个 Java 状态模型固定不变量
下面的模型不连接真实 Broker。它们把本篇最容易混淆的状态、位置或所有权变化压缩成确定输出;随后再用 RabbitMQ、Kafka 或 RocketMQ 的集成环境验证协议、持久化、重平衡和故障时序。
javac --release 17 -Xlint:all -Werror examples/backend-development/message-event/outbox-cdc-domain-events/OutboxCommitDemo.java examples/backend-development/message-event/outbox-cdc-domain-events/SchemaCompatibilityDemo.java
java -cp examples/backend-development/message-event/outbox-cdc-domain-events OutboxCommitDemo
java -cp examples/backend-development/message-event/outbox-cdc-domain-events SchemaCompatibilityDemotransactionCommitted=true orderVersion=7 outboxRows=1 eventId=evt-order-7 publishMayRepeat=true
readerVersion=2 acceptedSchemas=[1, 2] defaultsApplied=true breakingRenameRejected=true模型输出变化时,应定位是哪条业务不变量被修改,而不是只更新示例字符串。真实集成测试还要注入进程崩溃、响应丢失、Broker 切换和下游超时,因为这些窗口无法由纯内存模型模拟。
契约、权限与数据生命周期不能交给默认值
Broker 身份认证只证明客户端是谁,授权还要限制其能发布、订阅、声明和管理哪些资源。生产服务通常不应拥有删除 Topic、清空队列或修改保留策略的权限;消费服务不应能够伪造上游事件。跨环境凭证、测试 Topic 与生产 ACL 必须隔离,避免一次误配置把测试消费者接入生产消费组。
消息体、header、死信和 trace 会形成多份数据副本。个人信息和密钥应最小化、加密或令牌化,并定义 Broker 保留、死信保留、日志采样和备份删除策略。删除权请求不能只改业务主表,还要知道哪些事件属于不可变审计事实、哪些派生副本必须过期。压缩与加密发生顺序、密钥轮换和历史消息解密能力也需要纳入恢复演练。
用分层测试证明协议与业务同时成立
生产排障从一条业务事实反向追踪:数据库是否提交、Outbox 是否存在、Broker 是否有对应位置、哪个消费者获得过投递、幂等表记录什么、ack 或 offset 在哪里。若任何一段没有持久关联标识,团队就只能根据时间猜测。可解释的消息系统不要求每条消息全量 trace,但要求在异常样本上恢复完整因果链。
把设计决定写成可以被证伪的承诺
对每项承诺给出反例。例如“支付事件不丢”的反例是数据库已提交而 Outbox 不存在;“同订单有序”的反例是扩分区后同一 orderId 落入不同位置;“消费幂等”的反例是业务提交成功但幂等记录回滚;“十分钟恢复”的反例是下游安全吞吐减去生产流量后净排空速率不足。能构造反例,才能构造门禁和演练。
指标、日志与容量要围绕状态迁移
容量模型至少计算平均/峰值每秒消息数、平均/p99 字节、保留时间、副本倍数、重试放大、批处理密度、消费者处理时长与下游并发。恢复测试要在稳定生产流量仍存在时注入积压,验证净排空速率,而不是暂停生产后测一个理想峰值。
上线与演练以可恢复为验收标准
故障演练依次注入生产响应丢失、Broker 节点切换、消费者提交后崩溃、下游超时、永久坏消息和热点 key。每次演练都回答三件事:业务不变量是否保持、重复或延迟是否有持久证据、系统是否在容量预算内自动或人工收敛。
