数据库运维7 min read次阅读

CDC 实时数据同步实战(Debezium + Flink CDC)

为什么 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 即天然断点,重启从断点续传。

适合"读 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=incrementalschema_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 条实战底线

  1. binlog 必须 ROW 格式binlog_format=ROW),STATEMENT 无法精确回放,CDC 无解。
  2. binlog_row_image=FULL,否则 UPDATE/DELETE 缺前镜像,下游无法精确还原。
  3. 源库给 CDC 专用只读账号,只授 REPLICATION 权限,绝不用 root。
  4. 消费端一律按主键幂等写入(UPSERT),天然容错乱序与重复。
  5. binlog/WAL 保留期 ≥ 7 天,监控"剩余可续传空间",防 consumer 宕机后断点失效。
  6. PostgreSQL 逻辑复制槽会撑爆磁盘,必须监控 slot 消费进度
  7. 大 DDL 期间暂停 connector,并用 Schema Registry 做向后兼容演进。
  8. 快照阶段会短暂锁表,高并发库改用 incremental / schema_only 快照模式。
  9. Kafka 断点即续传点,offset 要持久化且不被自动重置到 earliest(否则重放全量)。
  10. CDC 是"同步"不是"备份",不替代备份;下游故障要能暂停+重放,别丢中间变更。
分享:

相关文章

评论区