数据库是系统的单一事实源,但搜索、数仓、缓存、下游服务都想要实时感知它的变化。变更数据捕获(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 让业务系统无需改造即可把每次变更变成可消费的事件,配合幂等写入与水位线监控,就能构建一套高可靠、可观测的实时数据同步底座,同时天然支撑发件箱模式的事件发布。
延伸阅读
- Spring 集成消息队列 — Kafka 消费与死信策略基础
- 数据库连接池、读写分离与分库分表 — 数据存储与同步目标端
- 分布式事务:Seata、TCC 与 Saga 模式详解 — 本地消息表与最终一致
- Java 监控诊断与可观测性 — 同步链路指标与告警
- Spring Boot 核心原理与自动配置 — Kafka 与数据源自动装配
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。