业务库的数据是实时的,而下游的数据仓库、缓存、搜索引擎往往是滞后的——传统 ETL 定时全量同步的延迟以"小时"计,且全量扫描对源库冲击大。CDC(Change Data Capture,变更数据捕获)改变了这一切:它从数据库的日志(MySQL binlog、PostgreSQL WAL、MongoDB oplog)里"偷听"每一次变更,以毫秒级延迟把增量数据流式推给下游。本指南深入 CDC 的底层原理(binlog/WAL 解析、位点管理、Exactly-Once)、主流工具(Canal、Debezium、Flink CDC)的架构与选型、与 Kafka/Flink 的组合链路、Schema 演进处理,以及生产级 CDC 管道的设计与避坑。
一、CDC 是什么:从"拉取"到"订阅"
1.1 传统同步 vs CDC
传统 ETL(拉取):
定时全量 SELECT → 传输 → 写入
· 延迟:小时级
· 成本:全表扫描,冲击源库
· 捕获不到"删了哪些行/改了什么"(只有当前状态)
CDC(订阅):
解析数据库日志 → 逐条变更事件 → 实时推送
· 延迟:毫秒级
· 成本:只读日志,几乎零影响源库
· 捕获完整变更(INSERT/UPDATE/DELETE + before/after)
ℹ️ 核心洞察:CDC 把数据库变成"事件源"——业务里的每一次写操作,都天然成为下游系统可以订阅的数据事件流。这是实时数仓、缓存同步、搜索索引、审计追踪的统一底座。
1.2 CDC 的典型场景
· 实时数仓:业务库 → CDC → Kafka → Flink → DWS/DW
· 缓存刷新:订单变更 → 实时刷新 Redis 缓存
· 搜索索引:商品变更 → 实时同步到 ES
· 微服务数据共享:订单库变更 → 同步给下游服务
· 审计与合规:记录所有变更(before/after)
· 数据迁移/双写过渡
二、CDC 的底层原理:binlog 与 WAL
2.1 MySQL binlog
binlog 三种格式:
STATEMENT:记录 SQL(节省空间,但非确定性函数有风险)
ROW :记录每行变更前后值(推荐,信息最全)
MIXED :混合,默认 statement 必要时转 row
CDC 必须用 ROW 格式:
ROW 事件包含 before_image / after_image
→ 才拿得到完整的字段变化
binlog 三种日志记录方式:
· 简单轮询读取日志文件
· 基于 binlog 位点(position)断点续传(Canal 传统方式)
· GTID(全局事务 ID):更可靠的断点恢复(推荐)
-- 确认 binlog 配置
SHOW VARIABLES LIKE 'binlog_format'; -- 需为 ROW
SHOW VARIABLES LIKE 'gtid_mode'; -- 建议 ON
SHOW VARIABLES LIKE 'log_bin'; -- 需 ON
2.2 PostgreSQL WAL / 逻辑复制
PG 的 CDC 走"逻辑复制":
· 物理复制:WAL 原始字节(用于流复制,对库不可读)
· 逻辑复制:WAL 解码为 SQL 语义的变更流(CDC 用)
· 通过 pgoutput / wal2json 插件输出 JSON 变更
PG CDC 特点:
· 发布/订阅模式(PUBLICATION/SUBSCRIPTION)
· 天然支持 schema 变化与过滤
· Debezium 等工具消费逻辑复制输出
2.3 位点(Offset)管理
CDC 必须记录"读到哪了"(offset/position/GTID):
· 进程重启 → 从上次位点继续,不丢不重
· 存到 Kafka offset / state store / DB
Exactly-Once 挑战:
· 源日志消费(offset)与下游写入(Kafka/DB)跨系统
→ 需"幂等 + 去重"或依赖 Flink checkpoint 的端到端一致
三、主流工具架构
3.1 Canal(阿里,MySQL → 消息队列)
架构:
Canal Instance 伪装成 MySQL 从库
→ 拉取 binlog
→ 解析为结构化变更事件(CanalEntry)
→ 推给下游(Kafka / RocketMQ / DB)
适用:
· 纯 MySQL → Kafka/MQ 的轻量管道
· 已有 Kafka 生态的团队
· 对 Debezium 生态陌生、想更贴近阿里实践
# canal.properties / instance.properties
canal.instance.master.address=mysql-master:3306
canal.instance.dbUsername=canal
canal.instance.dbPassword=***
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=orders_db\\..*
# 位点持久化
canal.instance.gtidon=false
canal.instance.master.position=
3.2 Debezium(Red Hat,跨库 CDC 框架)
架构:
Debezium Connector(Kafka Connect 插件)
→ 连接源库(MySQL/PG/MongoDB/Oracle/SQL Server...)
→ 解析日志 → 输出到 Kafka topic(每个表一个 topic)
→ Kafka Connect 管理位点(offset storage)
特性:
· 原生 Kafka 生态(Avro/JSON、schema registry 集成)
· 内置字段 flatten、过滤、列裁剪
· 支持 schema 变更事件(DDL)
· Exactly-Once 依赖 Kafka 事务 / Flink
// Debezium MySQL connector 配置(Kafka Connect REST API)
{
"name": "orders-cdc",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql-master",
"database.port": "3306",
"database.user": "debezium",
"database.password": "***",
"database.server.id": "5400",
"database.include.list": "orders_db",
"table.include.list": "orders_db.orders,orders_db.order_items",
"topic.prefix": "cdc",
"schema.history.internal.kafka.topic": "schema-history",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter"
}
}
// 消费到的事件(Kafka topic: cdc.orders_db.orders)
{
"op": "u", // c=create, u=update, d=delete
"before": {"id": 1001, "status": "PENDING", "amount": 500},
"after": {"id": 1001, "status": "PAID", "amount": 500},
"source": {"db": "orders_db", "table": "orders", "lsn": 123456}
}
3.3 Flink CDC(流式计算集成)
Flink CDC Connector 把"日志解析"内嵌进 Flink 作业:
· Flink SQL 直接同步:Source 是 CDC 流,Sink 是目标表
· 端到端 Exactly-Once(checkpoint + 幂等)
· 支持动态 schema 演进
· 一条 SQL 完成"实时入湖/入仓"
典型用法:
· 实时同步 MySQL → Kafka / Hudi / ClickHouse
· 双活/跨区同步
· 实时维表关联
-- Flink SQL:MySQL → Hudi 实时入湖
CREATE TABLE orders_cdc (
id INT, status STRING, amount DECIMAL(10,2),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql-master',
'database-name' = 'orders_db',
'table-name' = 'orders',
'server-id' = '5400-5404'
);
CREATE TABLE orders_ods (
id INT, status STRING, amount DECIMAL(10,2),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'hudi',
'path' = 's3://warehouse/ods/orders'
);
INSERT INTO orders_ods SELECT * FROM orders_cdc;
四、工具选型
| 工具 | 源支持 | 下游 | 特点 |
|---|---|---|---|
| Canal | MySQL | Kafka/MQ/RocketMQ | 阿里系,轻量 |
| Debezium | MySQL/PG/Mongo/… | Kafka Connect 生态 | 跨库、标准化 |
| Flink CDC | MySQL/PG/… | Flink 任意 Sink | SQL 化、Exactly-Once |
| Maxwell | MySQL | Kafka | 单表 JSON,极简 |
| Tapdata/流平台 | 多源 | 多目标 | 低代码产品 |
选型要点:
· 已上 Kafka Connect → Debezium
· 主要 MySQL + 阿里云栈 → Canal
· 要实时入湖/入仓/复杂计算 → Flink CDC
· 跨多种数据库统一 → Debezium / Flink CDC
· 团队熟悉度 + 运维成本
五、端到端管道架构
5.1 经典链路
业务库(MySQL/PG)
│ binlog/WAL
▼
CDC (Canal/Debezium/Flink CDC)
│ 变更事件流
▼
Kafka (每个表一个 topic / 分区按主键)
│ 多路分发(消费者组)
├──► Flink → 实时数仓 / ClickHouse / Hudi
├──► Redis 缓存刷新
├──► Elasticsearch 索引同步
└──► 审计/归档存储
关键设计:
· 分区键 = 主键(保证同一行事件有序)
· 下游消费幂等(upsert)
· schema 变化通过 registry 协调
5.2 增量 + 全量混合
初次建管道 = 全量快照 + 增量追平:
1. 先做全量初始同步(如已有备份/历史数据)
2. 启动 CDC 增量(位点对齐)
3. 合并去重 → 保证不丢不重
Flink CDC 内置此能力:启动时自动快照 + 增量衔接
六、Schema 演进处理
6.1 DDL 变更的影响
业务库加列/改类型 → CDC 事件结构变化 → 下游 schema 要跟上
三种处理:
· 统一 schema registry(Avro + Confluent Schema Registry)
· 宽松 JSON(新增字段容忍)
· Flink 动态列(重建表 schema)
Debezium 处理:
· DDL 事件也发到 schema-history topic
· 下游可从 schema history 重建表结构
6.2 兼容策略
· 只加列(可空/默认值):向后兼容,直接推进
· 改列类型:评估影响面(类型不兼容需重放)
· 删列:下游若引用则报错,需提前治理
· 建议:源库 DDL 走审批 + 通知 CDC 管道
七、生产级管道设计
7.1 可靠性
· 位点持久化(Kafka Connect offset / Flink checkpoint)
· 源库 binlog 保留时长(建议 24-72h,给管道容错窗口)
· 下游幂等(主键 upsert)
· 断点续跑:管道重启从位点继续,无重复无丢失
· 监控:事件延迟(lag)、处理吞吐、失败重试
7.2 事件保序
同一条数据(同主键)的事件必须有序:
· Kafka 分区按主键 hash → 同 key 同分区 → 有序
· 下游单分区消费或按 key 聚合处理
· 跨表(如 order + order_item)天然独立 topic,无需保序
7.3 反压与限流
· 源库 binlog 解析跟不上 → 事件积压 → 延迟上升
→ 扩容 CDC 消费者 / 增大 Kafka 分区
· 下游慢 → 背压传播
→ 用 Kafka 缓冲解耦 + 独立扩容消费组
· 大事务/批量变更 → 事件风暴
→ 下游幂等 + 限流(最大吞吐阈值)
八、监控与运维
8.1 关键指标
# CDC 健康指标
cdc_events_total{table} # 事件量
cdc_event_latency_seconds # 端到端延迟(源库→下游)
kafka_consumer_lag{group, topic} # 消费积压
cdc_parse_error_total # 解析错误
cdc_schema_version{table} # schema 版本
cdc_binlog_position # 位点进度(对比最新)
# 告警:lag 超阈值、parse error 持续、位点停止
8.2 日常运维
· 版本升级:先扩容新版本实例 → 灰度切流
· 位点重置:确认不丢数据才允许
· 表结构变更:走 DDL 审批并同步更新下游
· 故障演练:模拟 binlog 中断 / 消费组故障
九、落地清单与避坑
9.1 Checklist
□ binlog_format=ROW + GTID 开启(MySQL)
□ 授权 CDC 账号(REPLICATION SLAVE/CLIENT 等)
□ 选型:Kafka Connect/Debezium vs Flink CDC vs Canal
□ Kafka topic 设计(每表一 topic,主键分区)
□ 全量+增量衔接(初次建管道)
□ schema registry / 演进策略
□ 位点持久化 + 断点续跑验证
□ 下游幂等 upsert
□ 延迟/积压监控 + 告警
□ binlog 保留时长评估
9.2 常见坑
| 坑 | 现象 | 对策 |
|---|---|---|
| binlog 非 ROW | 无 before/after | 改 ROW + 换新管道 |
| 位点丢失 | 重复/丢失事件 | 持久化 offset |
| 下游不幂等 | 重复数据 | 主键 upsert |
| schema 变更崩 | 解析失败 | registry + DDL 审批 |
| binlog 过期 | 无法回追 | 保留 24-72h |
| 事件风暴 | 下游被冲垮 | 限流 + 缓冲 |
| 跨表保序误判 | 数据错乱 | 分区键 = 主键 |
9.3 反模式
· 用轮询 SELECT 代替 CDC(重复拉取、伤源库)→ 应走日志
· 下游直接依赖"事件顺序"做业务逻辑(跨表无序)
· 把 CDC 当无限准确(大事务丢细粒度)→ 配合幂等
· 忽略 schema 演进(加列即崩)
总结:CDC 决策表
| 环节 | 关键动作 |
|---|---|
| 原理 | binlog/WAL 解析 + 位点管理 |
| 工具 | Canal(MySQL/MQ)/ Debezium(Kafka)/ Flink CDC(SQL) |
| 链路 | 源库 → CDC → Kafka → Flink/缓存/搜索/数仓 |
| 可靠 | 位点持久化 + 幂等 + 断点续跑 |
| Schema | registry + DDL 审批 |
| 保序 | 分区键 = 主键 |
| 监控 | 延迟 / 积压 / 解析错误 |
CDC 重新定义了"数据库与其他系统的连接方式"——不再是定时拉取快照,而是实时订阅变更日志。它让数据层第一次拥有了统一的实时通道:缓存、搜索、数仓、审计共享同一条事件流。落地记住四件事:源库日志要开 ROW + 位点要持久化、下游必须幂等、schema 演进要提前治理、延迟与积压要可观测。把 CDC 管道建好,你的数据基础设施就从"批处理的世界"迈入了"流的世界"——实时不再是一句口号,而是默认能力。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。