Skip to content

CDC与Outbox:MySQL到MQ、ES、Redis的可靠同步与重放

CDC 和 Outbox 解决的是“事实源数据库已经提交后,怎样可靠地把变化传播到 MQ、ES、Redis、报表、搜索视图和其他服务”。它们不追求跨多个系统强一致,而是用可重试、可幂等、可重放、可对账的事件链路实现最终一致。

相关基础可继续阅读:可靠消息与Outbox缓存一致性MySQL与ES一致性MQ消费幂等

学习目标

学完本页要能回答:

  1. CDC、Outbox、本地消息表、事务消息分别解决什么问题。
  2. 为什么 MySQL 到 ES、Redis、MQ 不能靠普通双写保证一致。
  3. 捕获业务表 binlog 和捕获 Outbox 表有什么区别。
  4. CDC 链路里位点、顺序、重复、乱序、Schema 演进和回放怎么处理。
  5. ES 更新失败、缓存删除失败、消息消费失败分别怎么补偿。
  6. 为什么事件必须有 eventId、businessKey、version 和 schemaVersion。
  7. 线上数据不一致时按什么证据排查。

一、为什么需要 CDC 与 Outbox

商业系统里,数据库通常是事实源,其他系统是派生视图。

mermaid
flowchart TD
    A["MySQL事实源"] --> B["MQ业务事件"]
    A --> C["ES搜索视图"]
    A --> D["Redis缓存副本"]
    A --> E["报表和数仓"]
    A --> F["其他微服务本地视图"]

如果业务代码直接写多个系统:

mermaid
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 直接订阅业务表变更,比如 productorderasset

mermaid
flowchart TD
    A["业务更新product表"] --> B["MySQL写binlog"]
    B --> C["CDC读取变更"]
    C --> D["转换为商品变更事件"]
    D --> E["同步ES或缓存"]

优点是业务代码少,缺点是 CDC 消费者必须从行变更推断领域含义。比如商品价格变了、上下架状态变了、库存变了,是否都要更新同一份 ES 文档?是否都要通知下游?这些语义不一定能从单行变化里可靠推断。

3.2 捕获 Outbox 表

业务事务显式写入领域事件,再由 CDC 捕获 Outbox。

mermaid
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串联写入、投递、消费链路
statusPENDING、SENT、FAILED、DEAD
retry_count投递或处理次数

最关键的是 aggregate_version。ES、Redis、本地视图更新时要比较版本:新版本覆盖旧版本,旧版本晚到时丢弃或忽略。

五、MySQL 到 ES 同步链路

mermaid
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 兜底”。

mermaid
flowchart TD
    A["写请求提交MySQL"] --> B["事务提交后主动删除Redis"]
    B --> C{"删除是否成功"}
    C -- "成功" --> D["下次读回源新值"]
    C -- "失败" --> E["记录重试或告警"]
    A --> F["binlog CDC捕获变化"]
    F --> G["再次删除相关缓存Key"]
    G --> H["TTL最终兜底"]

为什么 CDC 不一定替代主动删缓存?

  1. CDC 有延迟,主动删除能缩短不一致窗口。
  2. CDC 可能按表变更触发,复杂缓存 key 需要业务映射。
  3. 缓存值可能是多表聚合,单表 CDC 不一定知道所有受影响 key。
  4. 最稳妥的是主动删除、CDC补偿、TTL兜底、监控告警组合。

七、顺序、重复和乱序

CDC/MQ 链路通常是至少一次语义。必须假设事件会重复、乱序、延迟。

mermaid
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 演进

事件一旦发出去,就要考虑历史消费者和历史消息。

兼容原则:

  1. 新增字段要给默认值,旧消费者可忽略。
  2. 不直接删除字段,先让消费者升级,再停止写旧字段,最后删除。
  3. 不随意改变字段含义和类型。
  4. 事件带 schema_version
  5. 重放历史事件时,消费者要能处理旧版本载荷。

如果事件 Schema 不治理,回放历史消息、重建 ES 索引、灾备恢复时就会失败。

九、位点、水位和重放

CDC 可靠性的核心证据是位点。

概念含义
binlog file/positionMySQL binlog 读取位置
GTID全局事务标识,便于主从切换和恢复
WAL LSNPostgreSQL 日志位置
OffsetMQ 分区消费位置
Watermark已确认同步到某个时间或版本

重放不是把消息随便再发一遍,而是按确定范围、确定版本、确定幂等策略重新应用。

重放流程:

mermaid
flowchart TD
    A["确定事实源范围"] --> B["确定起止位点或业务版本"]
    B --> C["暂停或限速相关消费者"]
    C --> D["按eventId和version幂等重放"]
    D --> E["记录成功、失败和死信"]
    E --> F["对账目标视图"]
    F --> G["恢复正常消费"]

十、商业场景

10.1 医疗资产同步 ES

资产表是事实源,ES 用于按医院、科室、设备名、状态快速检索。

设计:

  1. 资产变更事务写 asset 表和 outbox_event
  2. 事件包含 assetIdhospitalIdassetVersioneventType
  3. CDC 捕获 Outbox 投递 MQ。
  4. ES 消费者按 assetId 更新文档,版本低于 ES 当前版本时忽略。
  5. 失败进入死信,补偿任务按 MySQL 事实源重建文档。
  6. 每天对账 MySQL 最新版本和 ES 文档版本。

10.2 支付成功后通知、积分和搜索同步

支付成功不能被通知、积分、ES 阻塞。支付事务只更新支付流水、订单状态和 Outbox,后续异步传播。通知失败可以重试,积分失败可以补偿,ES 失败可以重建;支付事实不能因为通知失败而回滚。

10.3 缓存删除失败兜底

商品更新后业务代码删除 Redis 失败,CDC 捕获商品变更后再次删除相关 key。即使 CDC 也失败,TTL 最终过期,监控发现删除失败和 CDC Lag 后告警处理。

十一、JDK 8 Demo:版本幂等更新

下面 Demo 演示 ES/Redis 派生视图怎样避免旧事件覆盖新状态。

java
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);
    }
}

输出:

text
9:new

这个 Demo 说明两件事:eventId 处理重复,version 处理乱序。生产里要把消费日志和业务写入放到同一事务,ES 则使用外部版本或脚本条件更新表达类似语义。

十二、生产 Runbook

12.1 MySQL 已更新,ES 没更新

  1. 查 MySQL 事实源版本。
  2. 查 Outbox 是否有对应事件。
  3. 查 CDC 位点是否越过该事件。
  4. 查 MQ 是否收到事件、分区和 offset。
  5. 查 ES 消费日志、失败日志、死信。
  6. 查是否 Mapping 冲突、版本冲突或批量请求被拒绝。
  7. 用业务主键从 MySQL 重建 ES 文档。

12.2 CDC Lag 变大

  1. 查数据库 binlog/WAL 产生速度。
  2. 查 CDC 连接、解析线程、序列化和下游发送耗时。
  3. 查 MQ Broker 写入是否慢。
  4. 查单表大事务或大字段事件。
  5. 查消费者是否被下游 ES/Redis/DB 限速。
  6. 扩容前确认分区、业务 key 和下游容量。

12.3 重放后数据仍不一致

  1. 确认重放范围是否覆盖缺失事件。
  2. 确认消费者是否按版本忽略了旧事件。
  3. 确认 Schema 版本是否兼容。
  4. 确认目标视图是否有其他写入通道。
  5. 对比事实源、事件日志、消费日志和目标版本。

十三、面试标准回答

CDC 是从数据库日志捕获已提交变化,Outbox 是把业务事实和待传播事件放到同一个本地事务。生产里更推荐业务写 Outbox,CDC 只负责可靠搬运,避免从业务表变更里猜领域语义。CDC/Outbox 解决的是数据库提交后事件传播的可靠性,不是跨系统强一致;链路仍然是至少一次,所以消费者必须用 eventId 去重,用 businessKey 和 version 防止乱序覆盖,用 schemaVersion 处理事件演进。MySQL 到 ES、Redis、MQ 同步失败时,要查事实源版本、Outbox状态、CDC位点、MQ offset、消费日志、死信和目标视图版本,必要时按事实源重放或重建。

十四、本章小结

CDC/Outbox 的核心不是“订阅 binlog 很高级”,而是把数据传播做成可证明的工程闭环:事实源唯一、事件有语义、位点可恢复、消费可幂等、失败可重试、死信可处理、视图可重建、结果可对账。