最终一致、Outbox 与 Inbox:可靠更新与重复处理
订单事务已经提交,进程却在发送消息前退出,搜索索引便缺少这笔订单。把发送提前也有风险:消息已被消费,数据库事务随后回滚,索引中反而多出不存在的订单。
Outbox 把“业务变更”和“需要发布的事件”放进同一个数据库事务。后续发布可以重复执行;消费端再把事件去重与业务写放进另一笔本地事务。两个事务之间仍然存在延迟,进度、重试和对账负责让落后的数据继续更新。
先明确哪个结果需要同步
业务事实、事件与读模型
订单库保存订单当前状态,事件描述已发生的变更,搜索索引或统计表按照自己的查询需求建立投影。它们可能使用不同字段、索引和保留策略,不需要逐字节相等。
订单库:订单 42,版本 7,状态 PAID
└── 事件:OrderPaid,eventId E7,订单版本 7
├── 搜索投影:展示已支付,已应用版本 7
└── 通知服务:该支付事件的通知处理记录“最终一致”需要明确更新规则和完成条件。对搜索投影,可以要求停止新增写入且传输恢复后,每个订单追到源库的相应版本;对通知,则需要区分待发送、发送成功和永久失败。消息被移入死信队列,只是改变了存放位置,业务尚未因此完成。
最终一致本身没有承诺某个固定时间。应用应给出可接受延迟、超过延迟后的读取方式和失败处理,例如订单详情读取订单库,搜索页允许短暂延迟;支付回调的关键判断不能依赖尚未追平的搜索索引。一致性模型的关系见 CAP 与一致性。
事件字段各自解决什么问题
| 字段 | 用途 | 常见误用 |
|---|---|---|
| eventId | 标识同一个已发生事件,供重复检测 | 每次发布重试都生成新 ID |
| aggregateId | 标识订单等业务对象,作为路由或分组键 | 把租户信息留给消费者猜测 |
| aggregateVersion / sequence | 判断该对象事件前后关系 | 用全局时间戳代替对象版本 |
| eventType、schemaVersion | 选择事件解析器与兼容处理 | 用同名字段悄悄改变单位或含义 |
| payload | 快照或增量数据 | 消费者不知道是否可跳过旧版本 |
| occurredAt、trace 关联 | 业务展示与诊断 | 把机器时间排序当作提交顺序 |
eventId 在源事务中确定,重复发布沿用它。aggregateVersion 应与业务更新在同一事务中推进;数据库 sequence 的分配次序不能直接代表提交次序,见 时间与 ID。敏感字段按消费需求最小化,事件一旦进入多个日志或备份,后续删除和权限管理会更复杂。
业务提交后,怎样可靠地发布
两张表在一个事务内变化
示例用一个订单的累计金额与版本展示增量事件。实际写入结构是:
BEGIN;
UPDATE event_lab.orders
SET amount = amount + 10, version = 1
WHERE id = 1 AND version = 0;
-- 应用检查更新行数必须为 1。
INSERT INTO event_lab.outbox(id, version, delta)
VALUES ('E1', 1, 10);
COMMIT;任一 SQL 失败,整个事务回滚。提交成功后,业务和待发事件都能查到;发布程序暂时不运行只会留下可继续处理的 Outbox。Events.write 通过 JDBC 的同一个连接执行这些操作,并检查条件更新行数。
外层事务覆盖不了另一个独立连接的自动提交,更覆盖不了普通 HTTP 发布。使用 ORM 时要核对 Outbox Repository 是否参加同一个事务、异步线程是否切断了事务上下文。把事件先放入内存队列再提交,仍会遇到进程退出导致队列丢失。
轮询与 CDC
轮询发布器查询未发布记录,把消息发给 Broker,再记录发布结果。多个发布器可以按分片领取,或用短事务领取一批并保存 owner、租约和重试时间。长期持有数据库行锁等待网络确认会扩大锁等待;领取提交后再发送,则要允许领取者崩溃后的重新领取。
CDC 发布器从数据库日志捕获已提交的 Outbox 变化,减少反复扫描业务表的成本。它仍需要保存源位点、维护连接器状态和日志保留,并处理初始快照与重启重放。Debezium 的 Outbox Event Router按事件 ID、聚合 ID 与 payload 转换记录;聚合 ID 可作为消息键,但后续并行处理仍须保持所需顺序。
| 方式 | 持久进度 | 主要运行压力 |
|---|---|---|
| 轮询 Outbox | 发布状态、领取租约、重试时间 | 扫描索引、更新表、积压清理 |
| CDC | 日志位点、连接器偏移、快照状态 | 日志保留、连接器恢复、Schema 演进 |
选哪种方式都应保留“数据库已提交、发布尚未完成”的可查询状态。CDC 配置与运行实验属于连接器专题;下面运行的是实际数据库加 RabbitMQ 的受控发布,不把它称为 CDC。
Confirm 之后仍可能重复发送
发布成功的确认与 Outbox 状态更新不在同一个事务中:
源事务提交 → 发布消息 → Broker 确认 → 标记 Outbox 已发布
↑
此处进程退出,下次会再次发布先标记再发送会漏消息;先发送再标记会重复。可靠重试通常选择后者,并让下游正确处理重复。
RabbitMQ publisher confirm 说明 Broker 对发布做出的确认,consumer acknowledgement 则是消费者对投递的确认,两者方向与职责不同。启用 mandatory 时,还需要处理没有路由目标的 return;一个被 return 的消息仍可能得到发布确认,不能只等到 ack 就标记业务发送成功。确认机制
消息持久性还依赖队列类型、持久化配置、复制方式和实际确认语义。生产队列的选择与故障测试见 RabbitMQ。本实验使用自动删除的独占队列,专门验证确认、重投与消费事务,退出后不保留队列,不覆盖 Broker 重启持久性。
Inbox 怎样与业务更新一起提交
去重记录就是业务事务的一部分
消费者收到事件后,先尝试在 Inbox 插入 eventId,然后更新投影,最后提交。在提交成功之前,不向 Broker 确认完成。
收到 E1
└── 消费者数据库事务
├── INSERT Inbox(E1, payload fingerprint)
├── 按版本更新订单投影
└── COMMIT
└── ACK E1重复 eventId 由唯一键仲裁。示例使用 INSERT ... ON CONFLICT DO NOTHING,检查实际插入行数;重复时读取既有摘要,确认同 ID 携带的仍是同一内容。PostgreSQL INSERT定义了冲突处理与返回行为。
int inserted = Db.execute(connection,
"INSERT INTO event_lab.inbox(id,fingerprint) VALUES(?,?) " +
"ON CONFLICT(id) DO NOTHING", event.id(), event.fingerprint());
if (inserted == 0) {
// 核对既有摘要;一致才返回“已处理”,不一致抛出冲突。
}
// 新事件的投影写入与 Inbox 一起提交。两个消费者同时插入同一 ID 时,数据库约束参与并发仲裁。示例使用默认 Read Committed,后续查询使用新的语句快照;提高隔离级别可能带来需要重试的序列化失败,应按整个本地事务重试。事务隔离
不要先独立提交 Inbox,再执行可能失败的业务写。业务回滚之后,下一次重投会被“已处理”记录挡住。也不要捕获普通唯一约束异常后假定 PostgreSQL 事务还能继续提交其他语句;采用冲突子句、保存点或重启事务,必须与实际数据库行为一致。
提交以后、ACK 以前的窗口
业务提交成功,消费者在 ACK 前断开,Broker 会重新投递。新的消费者再次尝试 Inbox 插入,发现同一事件已经处理,便可以确认消息而不重复修改业务。
如果业务涉及远端短信或支付,本地 Inbox 无法把那次远端调用纳入本地提交。可以在本地事务再写下一跳 Outbox,由远端支持的幂等键和结果查询处理重复。数据库内“应用一次”的保证不能直接扩展成任意外部副作用“只发生一次”。
增量事件遇到缺口必须等待或补齐
快照事件包含某对象在版本 7 的完整投影所需状态,消费者已有版本 6 时可以直接应用 7,再忽略已被覆盖的旧快照。这个策略要求快照字段完整、版本可比较,并正确表达删除。
增量事件则不同。E1 增加 10,E2 再增加 20,如果先收到 E2,就不能仅把版本设为 2 并加 20,随后忽略 E1。示例通过条件更新要求 currentVersion = eventVersion - 1:
UPDATE event_lab.projection
SET amount = amount + :delta, version = :event_version
WHERE id = :aggregate_id AND version = :event_version - 1;更新零行时回滚整笔事务,包括刚插入的 Inbox。可以将未来版本保存在单独的待排序区,等缺失版本到达后继续;也可以从可信源取得完整快照后重建。待排序记录与“已应用 Inbox”应有清楚区别,否则重试会永久跳过未应用事件。
同一聚合分配到同一分区有助于维持输入顺序,但并行消费、重试旁路和死信重放可能再次打乱完成顺序。顺序要求应一直延伸到投影提交位置。
运行真实发布、重投和回滚实验
下载 分布式状态实验工程。Linux 宿主准备 Docker Engine、Compose v2、unzip;Docker daemon 权限由开发环境管理员配置。以下资源全部位于独立实验网络,没有暴露数据库或 Broker 端口:
unzip distributed-state-lab.zip
cd distributed-state
docker compose --profile events up -d --wait
docker compose ps
mkdir -p .m2
LAB_DIR="$(pwd -P)"
docker run --rm --network ds14-state-network \
--user "$(id -u):$(id -g)" -e HOME=/tmp -e MAVEN_CONFIG=/m2 \
--entrypoint mvn \
--mount "type=bind,src=$LAB_DIR,dst=/work" \
--mount "type=bind,src=$LAB_DIR/.m2,dst=/m2" \
--workdir /work maven:3.9.12-eclipse-temurin-25 \
-B -ntp -Dmaven.repo.local=/m2 -Duser.home=/tmp \
-Dtest=OutboxTest test环境为 PostgreSQL 18.6、RabbitMQ 4.3.5、Java 客户端 5.33.0、pgJDBC 42.7.13、Maven 3.9.12。Java 源码目标为 17,也可把 Maven 镜像换成 maven:3.9.12-eclipse-temurin-17。构建使用宿主 UID/GID 与可写缓存,不挂载 Docker socket;数据库使用普通 lab 角色,Broker 使用仅供该隔离实验的 lab 用户和固定演示密码。
三个服务应先 healthy,再执行测试。第一次下载失败时检查已批准的镜像仓库和 Maven 仓库,或导入可信离线镜像与依赖缓存;不跳过 Broker 用例,也不把连接失败改成条件跳过。
OutboxTest 有 6 个实际测试:源事务回滚、消费事务回滚、增量缺口后重试、同 ID 内容冲突、条件快照修复,以及 Broker 的两个重复窗口。最后一个测试执行两次已确认的发布,再对一次已提交业务但尚未 ACK 的投递发送 NACK/requeue,检查后续重投标志和数据库:
published=2 deliveries=3 applied=1 inbox=1 amount=10
Tests run: 6, Failures: 0, Errors: 0, Skipped: 0这次 NACK 主动触发真实 Broker 重排队,用于稳定复现提交后尚未确认的窗口;它没有杀死消费者操作系统进程。测试使用有期限的 basicGet 读取和显式 ACK;生产的订阅式消费与线程模型见 RabbitMQ Java 客户端指南。
工程只维护一个聚合,事件格式为固定受控字符串,帮助看清事务和版本变化。项目接入时应采用有明确 Schema 的序列化格式、输入大小限制、租户归属和错误隔离;不能直接把实验 parser 当通用消息入口。
延迟、坏消息与对账修复
从各段持久进度定位停滞
| 观察结果 | 优先查找的位置 | 恢复后的确认 |
|---|---|---|
| 源库已变,Outbox 没有对应事件 | 是否同连接同事务写 Outbox,是否存在绕过入口 | 同一业务版本能找到事件 |
| 未发布 Outbox 持续增加 | 领取者、重试时间、Broker return/nack、连接认证 | 最老待发年龄下降,发布确认推进 |
| Broker 有积压,Inbox 不增长 | 消费者连接、权限、数据库阻塞和处理错误 | 提交速率恢复,未确认数受控 |
| Inbox 增长但业务不变 | 是否提前提交 Inbox,实际更新行数是否检查 | 原缺失结果修复,反例不再漏写 |
| 部分聚合长期缺版本 | 键路由、重试队列、丢失事件或毒消息 | 缺失版本补齐,后续增量继续应用 |
发布计数和消费计数可能受重复影响;按 eventId 和聚合版本关联,才能判断哪些业务仍落后。平均延迟容易掩盖少量长期停滞的对象,应同时关注最老未完成事件、分位延迟和永久失败量。
重试使用有上限的指数退避与抖动,并限制并发。协议不兼容、必填字段缺失等错误不会因为立刻重试而消失,应保存原始事件、解析版本和失败原因,进入可修复的隔离流程。含个人信息的 payload 不应完整写入告警或普通日志。
对账从源事实与实际投影比较
对账独立读取业务库和投影,按对象版本及关键字段比较。它可以发现发布逻辑遗漏、过去的消费 bug、人工改表和恢复后位点错误;只查消息系统“积压为零”找不到这些差异。
修复与正常消费会并发。对账准备回填源快照 3 时,消费者可能已经推进到版本 4;直接覆盖会回退数据。条件更新既检查目标仍是观察时的版本,也要求源版本更新,并一起保存快照覆盖水位:
UPDATE event_lab.projection
SET amount = :source_amount, version = :source_version,
baseline_version = :source_version
WHERE id = :id
AND version = :observed_local_version
AND version < :source_version;更新零行后重新读取并比较,而不是无限重放旧修复。示例先把投影推进到 2,验证 observed=1 的旧观察被拒绝,再以新的观察值应用快照 3;这项观察值比较是保守的并发冲突检测。随后处理增量 E4 推进到 4,再验证快照 3 不能覆盖它。真实源快照还要配合明确的版本、删除标记和读取一致性,防止修复过程中拼出跨版本字段。
baseline_version=3 表明这份可信完整快照已经包含到版本 3 的变化。迟到 E3 即使尚未出现在 Inbox,也应记录其摘要并无副作用地确认;否则它会永远因版本缺口检查而重试。示例把 Inbox 插入、读取并锁定投影水位和增量应用放在同一事务,验证 E3 重投可结束、同 ID 内容冲突仍拒绝、E4 能继续应用。这个水位只来自可信完整快照,不能取“最大收到的事件版本”。
若投影完全丢失,通常先读取带版本的源快照初始化,再从相应位置处理后续变化;增量 Inbox 与新基线的关系需一起设计。只恢复投影旧备份而保留最新 Inbox,会让缺失的结果因“已处理”而无法重建。
保留期覆盖重试、重放与灾备
Outbox 已发布记录的清理取决于发布完成语义和历史重放需求。Inbox 的保留期至少覆盖可能再次收到旧事件的窗口,包括 Broker 保留期、人工死信重放、离线消费者与数据库备份恢复。一个简单的“保留七天”无法适用于所有系统。
无法长期保存逐事件记录时,可以使用业务操作唯一键、聚合版本或按连续序列保存水位,但必须证明被丢弃记录的重复仍会被拦住。增量处理存在缺口时,不能把最大已见序号当作连续完成水位。
Schema 升级先让消费者兼容新旧字段,再逐步切换生产者;新增可选字段通常比改变金额单位容易兼容。重放旧事件时使用对应解析方式,失败记录也保留其原版本。应用回退前检查旧代码是否理解已经进入 Broker 的新事件。
恢复吞吐要为积压留出余量
每秒新到 300 条、消费者稳定完成 500 条时,理想净排空速度是 200 条/秒。实际速度还受重试、数据库索引和热点对象限制。恢复时突然把所有积压并发放开,可能再次耗尽连接池;按下游可承受速率逐步增加,具体计算与控制见 背压与恢复。
完成实验后回收本项目服务:
docker compose --profile events downPostgreSQL 具名卷保留,Broker 实验队列随容器移除。确认数据不再需要后使用 docker compose --profile events down -v 删除该项目数据库卷;测试日志和缓存按当前排查需要保留,不能提交含运行密码或原始业务 payload 的材料。
权威资料与规范地址
- Debezium Outbox Event Router:https://debezium.io/documentation/reference/stable/transformations/outbox-event-router.html
- RabbitMQ 确认与重投:https://www.rabbitmq.com/docs/confirms
- PostgreSQL INSERT / ON CONFLICT:https://www.postgresql.org/docs/18/sql-insert.html
- PostgreSQL 事务隔离:https://www.postgresql.org/docs/18/transaction-iso.html
- RabbitMQ Java 客户端指南:https://www.rabbitmq.com/client-libraries/java-api-guide
