Flink CDC 实战
概述
CDC(Change Data Capture)把数据库的增删改操作实时同步到数仓/消息系统。Flink CDC 连接器直接读 binlog/WAL,支持"全量 + 增量"无缝切换,是实时数仓的核心采集组件。本文讲清 MySQL/PostgreSQL CDC、Debezium 集成与生产要点。
一、CDC 是什么
1.1 概念
数据库变更 → 变更日志(binlog/WAL)→ 捕获成流
插入 → +I(新增)
更新 → -U(撤销旧值)+U(新值)
删除 → -D(删除)| 技术 | 原理 |
|---|---|
| 基于查询 | 轮询对比(DataX 等) |
| 基于日志 | binlog/WAL(Flink CDC) |
| 对比 | 查询 | 日志 |
|---|---|---|
| 实时性 | 延迟 | 准实时 |
| 侵入 | 查库 | 无侵入 |
| 完整性 | 难捕捉删除 | 全操作 |
二、MySQL CDC
2.1 前提
| 条件 | 配置 |
|---|---|
| binlog 开启 | binlog_format=ROW |
| 权限 | SELECT、REPLICATION SLAVE 等 |
2.2 表定义
sql
CREATE TABLE cdc_orders (
id BIGINT PRIMARY KEY NOT ENFORCED,
user_id BIGINT,
amount DECIMAL(10, 2),
create_time TIMESTAMP(3)
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql',
'port' = '3306',
'username' = 'flink',
'password' = '***',
'database-name' = 'orderdb',
'table-name' = 'orders',
'scan.incremental.snapshot.enabled' = 'true'
);2.3 关键参数
| 参数 | 说明 |
|---|---|
database-name | 库名(支持正则) |
table-name | 表名(支持正则) |
scan.startup.mode | initial/earliest-offset/latest-offset |
scan.incremental.snapshot.enabled | 增量快照 |
scan.incremental.snapshot.chunk.size | 分块大小 |
connect.timeout | 连接超时 |
三、PostgreSQL CDC
3.1 前提
| 条件 | 配置 |
|---|---|
| 逻辑复制 | wal_level=logical |
| 发布 | CREATE PUBLICATION |
| 复制槽 | 自动创建 |
3.2 表定义
sql
CREATE TABLE pg_orders (...) WITH (
'connector' = 'postgres-cdc',
'hostname' = 'pg',
'database-name' = 'orderdb',
'schema-name' = 'public',
'table-name' = 'orders',
'username' = 'flink',
'password' = '***',
'slot.name' = 'flink_slot'
);3.3 复制槽注意
| 注意 | 说明 |
|---|---|
| 槽管理 | 唯一名称,避免冲突 |
| 积压 | 未消费日志堆积,注意清理 |
| 快照 | 先快照后逻辑复制 |
四、Debezium 集成
4.1 关系
| 对比 | Flink CDC | Debezium |
|---|---|---|
| 定位 | Flink 连接器 | 独立 CDC 工具 |
| 部署 | 内嵌于 Flink | 独立服务/Kafka Connect |
| 输出 | 直接 Flink 流 | Kafka 消息 |
| 关系 | 底层兼容 Debezium 格式 | 可作 Flink 数据源 |
4.2 集成方式
| 方式 | 说明 |
|---|---|
| Debezium → Kafka → Flink | 传统链路 |
| Flink CDC 直连 | 简化链路(推荐) |
Debezium 方式:
MySQL → Debezium → Kafka(Debezium 格式)→ Flink 消费
Flink CDC 方式:
MySQL → Flink CDC Source → 直接处理4.3 Debezium 格式
sql
'format' = 'debezium-json'
-- 消费 Debezium 写入 Kafka 的变更消息五、全量 + 增量同步
5.1 同步流程
initial 模式:
1. 全量快照(并行分块读取)
2. 快照期间记录 binlog 位点
3. 快照完成 → 从位点继续增量
4. 无缝衔接,不丢不重5.2 并行快照
| 特性 | 说明 |
|---|---|
| 分块 | 按主键切块并行读 |
| 动态分块 | 快照期间数据变化自动调整 |
| 支持 | 大表快速同步 |
5.3 位点管理
| 位点 | 说明 |
|---|---|
| binlog 位点 | 保存于 Flink 状态 |
| Checkpoint | 位点随快照持久化 |
| 恢复 | 从位点继续,不重不丢 |
六、同步到下游
6.1 常见链路
1. MySQL → Flink CDC → Kafka(消息化)
2. MySQL → Flink CDC → 数仓表(ODS)
3. MySQL → Flink CDC → MySQL/PG(同步)
4. MySQL → Flink CDC → Iceberg/Hudi(数据湖)6.2 Upsert 写数仓
sql
INSERT INTO dwd.ods_orders
SELECT * FROM cdc_orders;
-- 结果表需主键,CDC 变更按主键 upsert6.3 全量同步到 Kafka
sql
CREATE TABLE kafka_orders (...) WITH (
'connector' = 'kafka', ...
);
INSERT INTO kafka_orders SELECT * FROM cdc_orders;七、生产实践要点
| 要点 | 说明 |
|---|---|
| 开启 Checkpoint | 位点持久化 |
| 主键必配 | CDC 变更按主键 |
| 大表并行 | 增量快照 + 合理 chunk |
| 监控延迟 | 同步延迟与积压 |
| 权限最小化 | 专用账号 |
| Schema 变更 | 评估字段变更影响 |
| 坑 | 处理 |
|---|---|
| binlog 保留期短 | 调大 binlog 保留 |
| 无主键表 | 快照慢,加主键 |
| 事务冲突 | 优化事务与 Sink |
| 版本兼容 | 连接器与数据库版本匹配 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 快照慢 | 大表,开增量快照、调整 chunk |
| 增量不生效 | 位点/权限问题,检查 binlog 配置 |
| 数据重复 | Sink 未幂等,加主键 upsert |
| 延迟高 | 下游慢,扩容或优化 |
| 复制槽冲突 | PostgreSQL 槽名唯一 |