Flink CDC 实时入湖架构实战:MySQL 到 Apache Iceberg 的无锁同步、Row-Level Upsert 与 Schema 自动演进
在实时湖仓一体(Real-time Lakehouse)建设中,将业务核心 OLTP 数据库(MySQL / PostgreSQL)的变更数据实时同步至数据湖,是打通“交易系统”与“分析系统”的生命线。
然而,传统的 CDC 入湖方案普遍面临三大工程噩梦:
基于 Flink CDC 3.0(无锁分片读取 + 整库同步 Pipeline) 与 Apache Iceberg v2(行级 Upsert + Schema 自由演进),技术团队能够构建起一套**“全量无锁切片 + 增量流式追踪 + DDL 动态演进 + 零停机热升级”**的工业级实时入湖流水线。
本文深入剖析 Flink CDC 底层无锁切片算法、Iceberg Equality Delete 文件流转,并给出生产级 Flink CDC 3.0 整库同步与 Schema 演进实战。
一、底层核心:Flink CDC 3.0 无锁分片与增量平滑切换
Flink CDC 3.0 彻底抛弃了传统的全局锁机制,采用 增量快照算法(Incremental Snapshot Chunking):
+———————————————————————————–+
| 1. 全量快照阶段: 无锁分片并行读取 (Lock-free Chunk Parallel Reading) |
| – 将单张亿级大表按主键范围切分为 N 个 Chunk (如每 10 万行一个 Chunk) |
| – 多个 Task 并发读取 Chunk 快照,期间全程**不加任何全局读锁 (零锁表影响)** |
| – 记录每个 Chunk 的起始位点 (Low Watermark) 与结束位点 (High Watermark) |
+———————————————————————————–+
|
v (分片读完,精准位点对齐)
+———————————————————————————–+
| 2. 增量流式阶段: 单连接 Binlog 共享消费 (Shared Binlog Streaming) |
| – 自动聚合数十张分表,共享单一 Binlog 读取连接 (彻底消除连接风暴) |
| – 从历史所有 Chunk 的最高 Watermark 开始,平滑无缝切换为纯 Binlog 增量追踪 |
+———————————————————————————–+
|
v (Schema 动态变更捕获)
+———————————————————————————–+
| 3. Schema 自动演进阶段 (Online Schema Evolution) |
| – 实时捕获源端 `ALTER TABLE` DDL 事件 |
| – 自动向下游 Iceberg Catalog 发起 Schema 更新 (加列/类型拓宽) |
| – **作业全程不重启、流计算零中断!** |
+———————————————————————————–+
二、Iceberg v2 行级 Upsert 与 Delete 文件底层流转
CDC 数据包含大量的 UPDATE 和 DELETE 语义。在 Apache Iceberg v2 中,通过开启 write.upsert.enabled = true 实现高效的行级更新:
| INSERT (+I) | 直接写入新的 Data Parquet 数据文件。 |
| DELETE (-D) | 生成对应的 Equality Delete 文件,记录被删除行的主键值; 下游读取时自动过滤该主键对应的数据。 |
| UPDATE (-U / +U) | 先针对旧主键写一条 Equality Delete 标记删除, 再写入一条包含新列值的数据文件,完成逻辑覆盖。 |
三、生产级 Flink CDC 3.0 整库同步与 Schema 演进配置实战
在 Flink CDC 3.0 中,推荐使用声明式 Pipeline YAML 规范 替代繁琐的单表 SQL,一条配置即可完成数十张分库分表的整库入湖:
# ====================================================================
# Flink CDC 3.0 生产级整库入湖 Pipeline: MySQL -> Iceberg
# ====================================================================
source:
type: mysql
hostname: mysql-cluster-prod.internal
port: 3306
username: cdc_lakehouse_user
password: "${env:MYSQL_CDC_PASSWORD}"
tables: "trade_db.orders_.*, trade_db.order_detail_.*" # 支持分库分表正则匹配
server-id: 5400-5408
server-time-zone: "Asia/Shanghai"
# 🌟 开启全量无锁并发分片读取
scan.incremental.snapshot.enabled: true
scan.incremental.snapshot.chunk.size: 8096
scan.snapshot.fetch.size: 2048
sink:
type: iceberg
catalog-type: hive
uri: "thrift://hive-metastore:9083"
warehouse: "s3://lakehouse-production-warehouse/iceberg"
database-name: "lakehouse_trade"
# 🌟 核心参数: 开启 Row-level Upsert 与 Schema 自动演进
table.properties:
format-version: "2"
write.upsert.enabled: "true"
write.target-file-size-bytes: "268435456" # 256MB 单文件大小
write.metadata.previous-versions-max: "5"
pipeline:
name: "mysql_to_iceberg_trade_cdc_pipeline"
parallelism: 8
# 开启两阶段提交 Checkpoint,保证端到端精确一次
checkpoint.interval: 60s
checkpoint.timeout: 180s
# 🌟 开启 Schema 变更自动向下游同步演进
schema.change.behavior: evolve
对应的 Flink SQL 单表深度定制脚本(含时间隐藏分区)
如果需要进行自定义字段转换或提取隐藏分区,可以使用 Flink SQL API:
— 1. 开启容错与状态后端
SET 'execution.checkpointing.interval' = '60s';
SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';
SET 'state.backend.type' = 'rocksdb';
— 2. 注册 MySQL CDC 增量源表
CREATE TABLE source_mysql_orders (
order_id BIGINT,
buyer_id BIGINT,
pay_amount DECIMAL(18, 2),
order_status STRING,
create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql-prod.internal',
'port' = '3306',
'username' = 'cdc_reader',
'password' = '${secret}',
'database-name' = 'trade_db',
'table-name' = 'orders',
'scan.incremental.snapshot.enabled' = 'true', — 无锁分片
'scan.startup.mode' = 'initial' — 先全量快照,后增量
);
— 3. 注册 Iceberg 目标湖仓表 (支持 Upsert 与按天隐藏分区)
CREATE CATALOG iceberg_catalog WITH (
'type' = 'iceberg',
'catalog-type' = 'rest',
'uri' = 'http://iceberg-rest-catalog:8181/v1',
'warehouse' = 's3://lakehouse-prod/'
);
CREATE TABLE IF NOT EXISTS iceberg_catalog.trade_lake.dwd_orders_realtime (
order_id BIGINT,
buyer_id BIGINT,
pay_amount DECIMAL(18, 2),
order_status STRING,
create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
dt DATE,
PRIMARY KEY (order_id, dt) NOT ENFORCED
) PARTITIONED BY (dt)
WITH (
'format-version' = '2',
'write.upsert.enabled' = 'true'
);
— 4. 实时流式写入
INSERT INTO iceberg_catalog.trade_lake.dwd_orders_realtime
SELECT
order_id,
buyer_id,
pay_amount,
order_status,
create_time,
update_time,
CAST(create_time AS DATE) AS dt
FROM source_mysql_orders;
四、生产避坑与高可用运维铁律
在维护大规模 Flink CDC 入湖流水线时,必须牢记以下四项生产铁律:
+—————————————————————————————–+
| Flink CDC 入湖避坑清单 |
|—————————————————————————————–|
| 1. 源表主键(Primary Key)是 Upsert 语义的唯一生命线: |
| – MySQL 源表必须包含明确的主键,且 Iceberg 目标表必须声明相同的主键; |
| – 若源表无主键,CDC 的 `UPDATE` 会退化为单纯的追加写入,导致湖中出现翻倍的重复数据! |
| |
| 2. 小文件合并(Compaction)与 Checkpoint 间隔的平衡: |
| – Checkpoint 间隔不宜低于 30 秒(推荐 1~3 分钟),否则会生成海量几百 KB 的碎片文件; |
| – 必须在后台配套常态化的 `rewrite_data_files` 定时任务,消除 Equality Delete 堆积。 |
| |
| 3. Binlog 保留时效(`binlog_expire_logs_seconds`)必须 >= 7 天: |
| – 如果 Flink 集群在节假日发生故障并停机排查 2 天,而 MySQL 的 Binlog 仅保留 24 小时, |
| 位点被物理清理后作业将无法断点续传,只能被迫推倒全量重算! |
| |
| 4. Schema 演进必须严格遵循“向下兼容”原则: |
| – 允许的自动演进:增加列(`ADD COLUMN`)、字段类型拓宽(`INT -> BIGINT`); |
| – 严禁操作:物理删除列或缩小字段精度,涉及此类破坏性变更必须走离线割接流程。 |
+—————————————————————————————–+
通过将 Flink CDC 3.0 无锁分片读取、整库同步 Pipeline 与 Apache Iceberg v2 行级 Upsert 深度协同,数据团队能够将传统 T+1 的数仓抽取彻底升级为秒级低延迟的现代实时湖仓底座。


