为什么 ETL 不够用了
传统数据同步靠定时批量抽取:
-- 每天凌晨抽昨天的数据
INSERT INTO dw.orders SELECT * FROM oltp.orders
WHERE update_time >= '2026-09-06 00:00:00';
问题显而易见:T+1 延迟、全表扫描压力大、删除/更新难以精确捕获(只认 update_time 会漏改)。当下游是实时大屏、风控、搜索索引、缓存失效时,小时级延迟直接不可用。
CDC(Change Data Capture,变更数据捕获)干的事是:实时捕获数据库的增删改,毫秒级同步到下游。它不扫全表,只读取数据库自己的事务日志(MySQL binlog、PostgreSQL WAL、Oracle redo)。
一句话区分:ETL 是"问数据库要数据",CDC 是"听数据库讲它刚改了什么"。
两种主流实现路线
| 方案 | 代表 | 架构 | 适合场景 |
|---|---|---|---|
| 日志中间件型 | Debezium | 抓取组件 → Kafka → 消费端 | 多源汇聚、解耦、需重放 |
| 计算引擎内置型 | Flink CDC | Flink 直接读 binlog → 写出 | 带转换/聚合、流批一体 |
两者底层都依赖源库的 binlog/WAL,本质一样,差别在"要不要经过 Kafka 中转"。
路线一:Debezium(经 Kafka)
前置:源库开启 binlog
-- MySQL 必须开启 ROW 格式 binlog(STATEMENT 无法精确回放)
-- my.cnf
[mysqld]
server-id = 1
log_bin = mysql-bin
binlog_format = ROW # 关键!STATEMENT/MIXED 不行
binlog_row_image = FULL # 捕获前镜像,便于 DELETE/UPDATE
expire_logs_days = 7
# 给 Debezium 一个只读账号
CREATE USER 'debezium'@'%' IDENTIFIED BY 'Debezium#2026';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'debezium'@'%';
Debezium MySQL Connector 配置
{
"name": "mysql-orders-cdc",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "10.10.4.17",
"database.port": "3306",
"database.user": "debezium",
"database.password": "Debezium#2026",
"database.server.id": "184054",
"database.server.name": "mysql-prod",
"database.include.list": "orders",
"table.include.list": "orders.t_order",
"database.history.kafka.bootstrap.servers": "kafka:9092",
"database.history.kafka.topic": "schema-changes.orders"
}
}
启动后,每张表的增删改会实时写入 Kafka topic:
mysql-prod.orders.t_order # 主题名 = server.schema.table
变更事件结构(简化):
{
"payload": {
"op": "u", // c=insert u=update d=delete r=snapshot
"before": { "id": 1001, "status": "PAID" },
"after": { "id": 1001, "status": "SHIPPED" },
"source": { "ts_ms": 1725600000000, "file": "mysql-bin.000003", "pos": 1543 },
"ts_ms": 1725600000123
}
}
下游消费时按
op分支处理:c/u写最新值、d删键。Kafka 的 offset 即天然断点,重启从断点续传。
路线二:Flink CDC(免 Kafka 直读)
适合"读 binlog → 简单转换 → 写目标库"的链路,少一个组件:
// Flink SQL: MySQL → 实时宽表
CREATE TABLE orders (
id BIGINT, status STRING, amount DECIMAL(10,2),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '10.10.4.17',
'port' = '3306',
'username' = 'debezium',
'password' = 'Debezium#2026',
'database-name' = 'orders',
'table-name' = 't_order',
'scan.startup.mode' = 'initial' // initial=先全量快照再增量;latest=只增量
);
CREATE TABLE doris_orders (
id BIGINT, status STRING, amount DECIMAL(10,2),
PRIMARY KEY (id) NOT ENFORCED
) WITH ('connector' = 'doris', 'table-name' = 'orders_rt');
INSERT INTO doris_orders SELECT id, status, amount FROM orders;
scan.startup.mode='initial'让 Flink 先做一次一致性快照(全量),再无缝切到 binlog 增量,业务零停机。
常见坑与治理
1. 数据乱序与重复
Kafka 单 partition 内有序,跨 partition(table.include.list 多表默认多分区)会乱序。按主键做幂等写入(UPSERT / 目标库 PRIMARY KEY + ON DUPLICATE KEY UPDATE)即可容错。
2. binlog 过期被清理
expire_logs_days 太小,Debezium 宕机超过保留期后无法续传。监控 Seconds_Behind_Source 与 binlog 剩余量,保留期设为至少 7 天。
3. DDL 变更中断
加列/改类型会让 Debezium schema 失败。用 schema-changes topic 留痕,配合 Avro Schema Registry 做兼容演进;大 DDL 期间临时暂停 connector。
4. 大事务 / 长事务
一个事务改 100 万行,Debezium 会攒巨量事件。控制源库单事务行数,或开启 log.retention 防止 Kafka 撑爆。
5. 全量快照锁表
Debezium 默认快照阶段会对表加 FLUSH TABLES WITH READ LOCK 短暂锁。高并发库用 snapshot.mode=incremental 或 schema_only 规避。
PostgreSQL 注意事项
PG 不用 binlog 而用 WAL + logical replication slot:
-- 需 superuser 建发布
CREATE PUBLICATION debezium_pub FOR TABLE orders.t_order;
-- 并修改 postgresql.conf
-- wal_level = logical
逻辑复制槽会一直保留 WAL 直到被消费,consumer 挂了 WAL 会无限增长撑爆磁盘——必须监控 slot 的 pg_replication_slots.confirmed_flush_lsn 与磁盘。
10 条实战底线
- binlog 必须 ROW 格式(
binlog_format=ROW),STATEMENT 无法精确回放,CDC 无解。 binlog_row_image=FULL,否则 UPDATE/DELETE 缺前镜像,下游无法精确还原。- 源库给 CDC 专用只读账号,只授
REPLICATION权限,绝不用 root。 - 消费端一律按主键幂等写入(UPSERT),天然容错乱序与重复。
- binlog/WAL 保留期 ≥ 7 天,监控"剩余可续传空间",防 consumer 宕机后断点失效。
- PostgreSQL 逻辑复制槽会撑爆磁盘,必须监控 slot 消费进度。
- 大 DDL 期间暂停 connector,并用 Schema Registry 做向后兼容演进。
- 快照阶段会短暂锁表,高并发库改用
incremental/schema_only快照模式。 - Kafka 断点即续传点,offset 要持久化且不被自动重置到 earliest(否则重放全量)。
- CDC 是"同步"不是"备份",不替代备份;下游故障要能暂停+重放,别丢中间变更。

