CDC 数据同步:Debezium 与 Kafka 架构实战

深入变更数据捕获原理与 Debezium 架构,讲解 Connector 配置、同步到搜索与数仓的管线、Outbox 联动,以及位点容错、幂等写入与同步延迟水位线的生产实践

数据库是系统的单一事实源,但搜索、数仓、缓存、下游服务都想要实时感知它的变化。变更数据捕获(CDC)让「数据库一改,外部系统立刻跟着改」成为现实,而 Debezium 是这条链路中最成熟的落地框架。本文从 CDC 原理、Debezium 架构、Connector 配置到搜索与数仓同步,给出完整的工程实现。

前置基础可先阅读 Spring 集成消息队列 与 数据库连接池、读写分离与分库分表。

1. CDC 概念与分类

1.1 什么是变更数据捕获

CDC 指捕获数据库中的增删改变化,并把它们作为有序事件流对外发布。与批量同步相比,CDC 是增量、近实时的,且对业务代码零侵入。

1.2 实现方式对比

方案原理侵入性延迟缺点
应用双写业务代码同时写库和外部系统高低不一致风险
定时轮询按时间戳/版本扫描变更表低中延迟与压力
触发器数据库触发器写日志表中低影响写性能
日志解析解析 binlog/WAL无极低需维护解析组件

日志解析是当前主流,Debezium 正是基于这一思路的实现。

2. Debezium 架构

2.1 核心组件

Debezium 以 Kafka Connect 插件形态运行,由三类角色组成:

数据库(MySQL binlog)
   │ 读取 binlog
   ▼
Debezium Source Connector(Kafka Connect Worker 内运行)
   │ 反序列化 + 快照 + 增量
   ▼
Kafka Topic(按表组织)
   │ 下游消费
   ▼
目标端(Elasticsearch / 数仓 / 业务服务)

2.2 快照与增量捕获

Connector 启动时先做全量快照,再无缝切换到增量 binlog 捕获。快照与增量的衔接由 Debezium 内部的 offset 机制保证,业务无感知。

时间线:
[全量快照快照段] → [增量 binlog 流] → [持续增量]
        offset 记录快照完成位置,增量从该位置继续

2.3 Topic 组织与事件结构

Debezium 默认按 topic.prefix.库名.表名 命名 topic,每条消息是包含 before、after、op、source 的包裹结构:

{
  "before": null,
  "after": {
    "id": 1001,
    "user_id": 88,
    "total": 99.90,
    "status": "PAID"
  },
  "op": "c",
  "ts_ms": 1727000000000,
  "source": { "db": "shop", "table": "orders", "ts_ms": 1726999999000 }
}

op 取值:c 创建、u 更新、d 删除、r 快照读取。source 中的 ts_ms 是数据库提交时间,正是计算同步延迟的关键字段。

3. Connector 配置

3.1 MySQL Connector 配置

# register-mysql.json
{
  "name": "orders-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql-primary",
    "database.port": 3306,
    "database.user": "debezium",
    "database.password": "changeit",
    "database.server.id": "5401",
    "database.include.list": "shop",
    "table.include.list": "shop.orders,shop.order_items",
    "database.history.kafka.topic": "schema-history.orders",
    "topic.prefix": "shop",
    "include.schema.changes": "true",
    "snapshot.mode": "initial"
  }
}

3.2 表过滤与字段转换

用正则控制同步范围,用单消息转换(SMT)重塑消息格式:

transforms: unwrap
transforms.unwrap.type: io.debezium.transforms.ExtractNewRecordState
# 去掉 Debezium 的复杂包裹结构,直接输出业务字段

3.3 快照模式选择

snapshot.mode行为适用场景
initial先快照后增量首次上线
when_needed无 offset 时快照断点续跑
schema_only只同步表结构已有增量起点
never不做快照存量已同步

3.4 数据类型与精度处理

MySQL 的 DECIMAL、时间、二进制类型需要显式映射,否则精度丢失:

decimal.handling.mode: precise      # 用字符串传输 DECIMAL,避免浮点误差
time.precision.mode: adaptive_time_microseconds
binary.handling.mode: bytes         # BINARY/VARBINARY 用字节数组

另外,database.server.id 必须唯一,同一实例下多个 Connector 的 server id 不能相同;主从切换场景要配置 gtid.new.channel.position 或开启 GTID 模式以自动定位位点。

4. 同步到搜索与数仓

4.1 同步到 Elasticsearch

订单表变化实时驱动搜索索引更新,是 CDC 最常见的落地场景:

@Component
public class OrderIndexSink {

    private final ElasticsearchClient esClient;

    @KafkaListener(topics = "shop.shop.orders", groupId = "order-indexer")
    public void onOrderChange(ConsumerRecord<String, String> record) {
        String key = record.key();                 // 主键
        String payload = record.value();           // 变更后的完整行
        OrderDocument doc = JsonUtils.fromJson(payload, OrderDocument.class);
        if (doc.isDeleted()) {
            esClient.delete(d -> d.index("orders").id(key));
        } else {
            esClient.index(d -> d.index("orders").id(key).document(doc));
        }
    }
}

也可以直接用 Kafka Connect 的 Elasticsearch Sink Connector,无需写代码:

# register-es-sink.json
{
  "name": "orders-es-sink",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "tasks.max": "3",
    "topics": "shop.shop.orders",
    "connection.url": "http://es:9200",
    "key.ignore": "false",
    "schema.ignore": "true",
    "behavior.on.null.values": "delete"
  }
}

4.2 数仓与数据湖

数仓同步通常把 CDC 流写入 ODS 层,配合 upsert 语义保持明细一致:

// 数仓 ODS 层:按主键 upsert,保留变更流水
public void upsertToWarehouse(OrderChange change) {
    warehouseDao.upsert("ods_orders", change.getKey(), change.getPayload());
    changeDao.append("ods_orders_changelog", change.getKey(),
            change.getOperation(), change.getTs());
}

4.3 宽表构建

订单与订单项拆分在多个表,通过 CDC 流按主键 join 出宽表供报表查询,避免在线库的 join 压力:

orders 表 CDC ─┐
                ├──> 按 order_id 聚合 → 宽表 ods_order_wide
order_items 表 ─┘

宽表构建要注意删除语义:父表删除时,聚合宽表对应的行也要同步清理,不能只做追加式 upsert。

4.4 同步一致性校验

CDC 只是增量管道,存量数据是否与源一致仍需周期校验。生产上普遍采用「抽样比对」:定期对目标端与源库按主键分片抽样,比对关键字段的哈希:

-- 源库侧计算分片哈希
SELECT id, MD5(CONCAT_WS('#', id, user_id, total, status)) AS row_hash
  FROM shop.orders
 WHERE id BETWEEN ? AND ?;

-- 目标端用同样函数计算后对比,不一致即触发告警

校验频率视数据重要性而定,核心业务表建议每小时做一次抽样一致性校验。

5. Outbox 联动

5.1 Outbox 事件路由器

Debezium 官方 SMT outbox-event-router 专门处理发件箱表,把数据库变更直接转成规范事件,实现「事务日志监听 + 发件箱」的黄金组合:

transforms: outbox
transforms.outbox.type: io.debezium.transforms.outbox.EventRouter
transforms.outbox.table.fields.additional.placement:
  type: event_type:header:eventType
transforms.outbox.route.by.field: event_type

5.2 与发件箱模式配合

业务方法写业务表 + 发件箱表(同一本地事务)
  → binlog 捕获发件箱 INSERT
  → outbox SMT 重塑为事件消息
  → Kafka 按事件类型路由到各 topic
  → 下游幂等消费

这套链路让「业务提交」与「事件发布」解耦,且发布动作完全由数据库日志驱动,延迟达到毫秒级。发件箱表设计要点见发件箱模式专题。

6. 容错与水位线

6.1 位点与故障恢复

Debezium 通过 Kafka Connect 的 offset 机制记录已消费的 binlog 位点。Connector 重启后从位点续跑,不丢不重。schema history 存在独立 topic,保证表结构变更也能恢复。

offset 状态:
  { connector: orders-connector,
    file: mysql-bin.000023,
    pos: 1048576,
    snapshot: false }

6.2 目标端幂等写入

CDC 至少投递一次,目标端必须幂等。写 Elasticsearch 按主键 index 天然幂等;写数仓用 upsert;写 Kafka 之外的系统要带业务主键做去重。

6.3 延迟监控水位线

同步延迟是 CDC 管线的核心健康指标,用「事件产生时间到目标端处理时间」的差值度量:

// 在 Sink 中记录端到端延迟
@KafkaListener(topics = "shop.shop.orders")
public void onOrderChange(ConsumerRecord<String, String> record) {
    long sourceTs = parseSourceTs(record.headers().lastHeader("ts"));
    long lag = System.currentTimeMillis() - sourceTs;
    lagGauge.set(lag);                      // 暴露为指标供告警
}

7. 生产实践

7.1 常见坑

坑现象对策
DDL 变更目标表结构与源不一致订阅 schema 变化 topic 并联动迁移
大事务binlog 事件洪峰调整批次与分区,增加缓冲
权限不足读不到 binlog授予 REPLICATION SLAVE 权限
位点漂移重复或丢失备份 offset,定期做一致性校验

7.2 监控指标

核心指标:
  - connector 任务状态(RUNNING / FAILED)
  - 每秒事件数(records_consumed)
  - binlog 剩余延迟(毫秒)
  - 目标端写入失败率
  - 端到端同步延迟水位线

7.3 双写 vs CDC 选型

维度应用双写CDC
侵入性业务代码改造无侵入
一致性双写难保证由数据库日志保证
实时性高高
复杂度低需运维 Connect 集群
适用场景简单旁路事件搜索、数仓、审计、事件流

8. 总结

主题核心要点
CDC 原理日志解析最优,binlog 驱动零侵入
Debezium 架构Source Connector + Kafka Connect + offset
Connector 配置表过滤、SMT 重塑、快照模式选择
目标同步搜索按主键索引、数仓 upsert、宽表聚合
Outbox 联动发件箱 SMT 把发布交给 binlog
容错水位线位点续跑、幂等写、延迟监控

CDC 把「数据产生方」与「数据消费方」之间的同步从定时批量进化为事件驱动的实时流。Debezium 让业务系统无需改造即可把每次变更变成可消费的事件,配合幂等写入与水位线监控,就能构建一套高可靠、可观测的实时数据同步底座,同时天然支撑发件箱模式的事件发布。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java-enterprise」更多文章

  1. JPMS 模块化:module-info 与 JLink 精简运行时
  2. 可观测性工程:Micrometer 指标模型与 OTLP 导出
  3. 多租户 SaaS 架构:隔离模型与租户上下文传递