CDC与Outbox:MySQL到MQ、ES、Redis的可靠同步与重放
CDC 和 Outbox 解决的是“事实源数据库已经提交后,怎样可靠地把变化传播到 MQ、ES、Redis、报表、搜索视图和其他服务”。它们不追求跨多个系统强一致,而是用可重试、可幂等、可重放、可对账的事件链路实现最终一致。
相关基础可继续阅读:可靠消息与Outbox、缓存一致性、MySQL与ES一致性、MQ消费幂等。
学习目标
学完本页要能回答:
- CDC、Outbox、本地消息表、事务消息分别解决什么问题。
- 为什么 MySQL 到 ES、Redis、MQ 不能靠普通双写保证一致。
- 捕获业务表 binlog 和捕获 Outbox 表有什么区别。
- CDC 链路里位点、顺序、重复、乱序、Schema 演进和回放怎么处理。
- ES 更新失败、缓存删除失败、消息消费失败分别怎么补偿。
- 为什么事件必须有 eventId、businessKey、version 和 schemaVersion。
- 线上数据不一致时按什么证据排查。
一、为什么需要 CDC 与 Outbox
商业系统里,数据库通常是事实源,其他系统是派生视图。
flowchart TD
A["MySQL事实源"] --> B["MQ业务事件"]
A --> C["ES搜索视图"]
A --> D["Redis缓存副本"]
A --> E["报表和数仓"]
A --> F["其他微服务本地视图"]如果业务代码直接写多个系统:
flowchart TD
A["更新MySQL成功"] --> B["更新ES"]
B --> C{"ES是否成功"}
C -- "失败" --> D["MySQL新值存在"]
D --> E["ES仍是旧值"]
C -- "超时" --> F["不知道ES是否更新"]问题在于 MySQL、ES、Redis、MQ 是不同提交边界。普通本地事务无法同时覆盖它们,一边成功一边失败就会产生分叉。CDC/Outbox 的思路是:先把业务事实提交到单一事实源,再用可靠事件链路传播变化。
二、概念区别
| 名称 | 核心做法 | 解决什么 | 不解决什么 |
|---|---|---|---|
| 本地消息表 | 业务表和消息表同库同事务提交 | 业务提交后一定有待发送消息 | 消费端仍可能重复、失败 |
| Outbox | 本地消息表的领域事件化模式 | 明确事件类型、版本、状态和重放 | 不保证消费者业务成功 |
| CDC | 从数据库变更日志捕获已提交变化 | 避免业务线程同步投递消息 | 位点、乱序、Schema、幂等要自己治理 |
| 事务消息 | Broker 半消息、本地事务、回查 | 本地事务与消息可见性一致 | 下游消费仍需幂等和补偿 |
| 双写 | 业务代码依次写 DB 和外部系统 | 简单直接 | 崩溃窗口大,不适合关键一致性 |
一句话:Outbox 让“事件意图”和业务事实同事务保存;CDC 负责把已提交变化可靠搬运出去。
三、两种 CDC 模式
3.1 捕获业务表
CDC 直接订阅业务表变更,比如 product、order、asset。
flowchart TD
A["业务更新product表"] --> B["MySQL写binlog"]
B --> C["CDC读取变更"]
C --> D["转换为商品变更事件"]
D --> E["同步ES或缓存"]优点是业务代码少,缺点是 CDC 消费者必须从行变更推断领域含义。比如商品价格变了、上下架状态变了、库存变了,是否都要更新同一份 ES 文档?是否都要通知下游?这些语义不一定能从单行变化里可靠推断。
3.2 捕获 Outbox 表
业务事务显式写入领域事件,再由 CDC 捕获 Outbox。
flowchart TD
A["本地事务写订单"] --> B["同事务写outbox事件"]
B --> C["MySQL提交并写binlog"]
C --> D["CDC捕获outbox行"]
D --> E["投递MQ或同步ES"]
E --> F["消费者幂等处理"]优点是语义清晰:事件名、业务主键、版本、载荷、重放策略都由业务定义。生产里更推荐“业务写 Outbox,CDC 搬运 Outbox”,尤其是订单、支付、资产变更、权限变更、搜索同步这类需要可审计和可重放的场景。
四、事件模型怎么设计
事件不能只有一个 JSON 字符串。至少要有:
| 字段 | 作用 |
|---|---|
event_id | 全局唯一事件ID,消费者去重 |
event_type | 事件类型,例如 OrderPaid |
aggregate_type | 聚合类型,例如 Order |
aggregate_id | 业务主键,例如订单号 |
aggregate_version | 业务版本,防止旧事件覆盖新状态 |
schema_version | 载荷结构版本,支持兼容演进 |
payload | 事件内容 |
created_at | 事件产生时间 |
trace_id | 串联写入、投递、消费链路 |
status | PENDING、SENT、FAILED、DEAD |
retry_count | 投递或处理次数 |
最关键的是 aggregate_version。ES、Redis、本地视图更新时要比较版本:新版本覆盖旧版本,旧版本晚到时丢弃或忽略。
五、MySQL 到 ES 同步链路
flowchart TD
A["MySQL提交商品version=8"] --> B["写Outbox ProductChanged v8"]
B --> C["CDC捕获事件"]
C --> D["投递Kafka或RocketMQ"]
D --> E["ES同步消费者"]
E --> F["按productId和version幂等更新"]
F --> G["失败重试或死信"]
G --> H["对账任务比较MySQL和ES版本"]ES 更新失败怎么办?
| 场景 | 处理 |
|---|---|
| 网络超时 | 不确认成功,按 eventId 重试 |
| ES 429 | 限速、退避、批量大小调小 |
| Mapping 冲突 | 进入死信,修正字段或新建索引重放 |
| 旧事件晚到 | 比较版本,旧版本不覆盖新文档 |
| 消费者宕机 | 从 MQ Offset 或 CDC 位点恢复 |
| 长期不一致 | 按 MySQL 事实源重建索引或补偿 |
不要把 ES 当事实源。ES 是可重建搜索视图,订单、资产、商品最终状态要以 MySQL 或业务事实服务为准。
六、MySQL 到 Redis 缓存同步
缓存同步更常见的组合是“业务提交后主动删缓存 + CDC 兜底”。
flowchart TD
A["写请求提交MySQL"] --> B["事务提交后主动删除Redis"]
B --> C{"删除是否成功"}
C -- "成功" --> D["下次读回源新值"]
C -- "失败" --> E["记录重试或告警"]
A --> F["binlog CDC捕获变化"]
F --> G["再次删除相关缓存Key"]
G --> H["TTL最终兜底"]为什么 CDC 不一定替代主动删缓存?
- CDC 有延迟,主动删除能缩短不一致窗口。
- CDC 可能按表变更触发,复杂缓存 key 需要业务映射。
- 缓存值可能是多表聚合,单表 CDC 不一定知道所有受影响 key。
- 最稳妥的是主动删除、CDC补偿、TTL兜底、监控告警组合。
七、顺序、重复和乱序
CDC/MQ 链路通常是至少一次语义。必须假设事件会重复、乱序、延迟。
flowchart TD
A["事件v8先产生"] --> B["事件v9后产生"]
B --> C["v9先到消费者"]
C --> D["ES更新到version=9"]
A --> E["v8晚到"]
E --> F["比较version发现旧事件"]
F --> G["忽略v8"]治理规则:
| 问题 | 解决 |
|---|---|
| 重复投递 | 消费日志表用 eventId 去重 |
| 同一聚合乱序 | aggregate_id 做分区 key,或消费者按版本幂等 |
| 跨聚合顺序 | 通常不保证,业务不要依赖全局顺序 |
| 旧事件覆盖新状态 | 写入目标时比较 aggregate_version |
| 消费失败 | 重试、死信、补偿、对账 |
八、Schema 演进
事件一旦发出去,就要考虑历史消费者和历史消息。
兼容原则:
- 新增字段要给默认值,旧消费者可忽略。
- 不直接删除字段,先让消费者升级,再停止写旧字段,最后删除。
- 不随意改变字段含义和类型。
- 事件带
schema_version。 - 重放历史事件时,消费者要能处理旧版本载荷。
如果事件 Schema 不治理,回放历史消息、重建 ES 索引、灾备恢复时就会失败。
九、位点、水位和重放
CDC 可靠性的核心证据是位点。
| 概念 | 含义 |
|---|---|
| binlog file/position | MySQL binlog 读取位置 |
| GTID | 全局事务标识,便于主从切换和恢复 |
| WAL LSN | PostgreSQL 日志位置 |
| Offset | MQ 分区消费位置 |
| Watermark | 已确认同步到某个时间或版本 |
重放不是把消息随便再发一遍,而是按确定范围、确定版本、确定幂等策略重新应用。
重放流程:
flowchart TD
A["确定事实源范围"] --> B["确定起止位点或业务版本"]
B --> C["暂停或限速相关消费者"]
C --> D["按eventId和version幂等重放"]
D --> E["记录成功、失败和死信"]
E --> F["对账目标视图"]
F --> G["恢复正常消费"]十、商业场景
10.1 医疗资产同步 ES
资产表是事实源,ES 用于按医院、科室、设备名、状态快速检索。
设计:
- 资产变更事务写
asset表和outbox_event。 - 事件包含
assetId、hospitalId、assetVersion、eventType。 - CDC 捕获 Outbox 投递 MQ。
- ES 消费者按
assetId更新文档,版本低于 ES 当前版本时忽略。 - 失败进入死信,补偿任务按 MySQL 事实源重建文档。
- 每天对账 MySQL 最新版本和 ES 文档版本。
10.2 支付成功后通知、积分和搜索同步
支付成功不能被通知、积分、ES 阻塞。支付事务只更新支付流水、订单状态和 Outbox,后续异步传播。通知失败可以重试,积分失败可以补偿,ES 失败可以重建;支付事实不能因为通知失败而回滚。
10.3 缓存删除失败兜底
商品更新后业务代码删除 Redis 失败,CDC 捕获商品变更后再次删除相关 key。即使 CDC 也失败,TTL 最终过期,监控发现删除失败和 CDC Lag 后告警处理。
十一、JDK 8 Demo:版本幂等更新
下面 Demo 演示 ES/Redis 派生视图怎样避免旧事件覆盖新状态。
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
public class VersionedEventApplyDemo {
static final class Event {
final String eventId;
final String key;
final long version;
final String value;
Event(String eventId, String key, long version, String value) {
this.eventId = eventId;
this.key = key;
this.version = version;
this.value = value;
}
}
static final class View {
final long version;
final String value;
View(long version, String value) {
this.version = version;
this.value = value;
}
}
static final class Consumer {
private final Map<String, Boolean> consumed =
new ConcurrentHashMap<String, Boolean>();
private final Map<String, View> view =
new ConcurrentHashMap<String, View>();
void onEvent(Event event) {
if (consumed.putIfAbsent(event.eventId, Boolean.TRUE) != null) {
return;
}
View old = view.get(event.key);
if (old == null || event.version > old.version) {
view.put(event.key, new View(event.version, event.value));
}
}
View get(String key) {
return view.get(key);
}
}
public static void main(String[] args) {
Consumer consumer = new Consumer();
consumer.onEvent(new Event("e9", "asset-1", 9, "new"));
consumer.onEvent(new Event("e8", "asset-1", 8, "old"));
consumer.onEvent(new Event("e9", "asset-1", 9, "new"));
View view = consumer.get("asset-1");
System.out.println(view.version + ":" + view.value);
}
}输出:
9:new这个 Demo 说明两件事:eventId 处理重复,version 处理乱序。生产里要把消费日志和业务写入放到同一事务,ES 则使用外部版本或脚本条件更新表达类似语义。
十二、生产 Runbook
12.1 MySQL 已更新,ES 没更新
- 查 MySQL 事实源版本。
- 查 Outbox 是否有对应事件。
- 查 CDC 位点是否越过该事件。
- 查 MQ 是否收到事件、分区和 offset。
- 查 ES 消费日志、失败日志、死信。
- 查是否 Mapping 冲突、版本冲突或批量请求被拒绝。
- 用业务主键从 MySQL 重建 ES 文档。
12.2 CDC Lag 变大
- 查数据库 binlog/WAL 产生速度。
- 查 CDC 连接、解析线程、序列化和下游发送耗时。
- 查 MQ Broker 写入是否慢。
- 查单表大事务或大字段事件。
- 查消费者是否被下游 ES/Redis/DB 限速。
- 扩容前确认分区、业务 key 和下游容量。
12.3 重放后数据仍不一致
- 确认重放范围是否覆盖缺失事件。
- 确认消费者是否按版本忽略了旧事件。
- 确认 Schema 版本是否兼容。
- 确认目标视图是否有其他写入通道。
- 对比事实源、事件日志、消费日志和目标版本。
十三、面试标准回答
CDC 是从数据库日志捕获已提交变化,Outbox 是把业务事实和待传播事件放到同一个本地事务。生产里更推荐业务写 Outbox,CDC 只负责可靠搬运,避免从业务表变更里猜领域语义。CDC/Outbox 解决的是数据库提交后事件传播的可靠性,不是跨系统强一致;链路仍然是至少一次,所以消费者必须用 eventId 去重,用 businessKey 和 version 防止乱序覆盖,用 schemaVersion 处理事件演进。MySQL 到 ES、Redis、MQ 同步失败时,要查事实源版本、Outbox状态、CDC位点、MQ offset、消费日志、死信和目标视图版本,必要时按事实源重放或重建。
十四、本章小结
CDC/Outbox 的核心不是“订阅 binlog 很高级”,而是把数据传播做成可证明的工程闭环:事实源唯一、事件有语义、位点可恢复、消费可幂等、失败可重试、死信可处理、视图可重建、结果可对账。
