Flink CDC 同步 MySQL 到 Iceberg 太慢?我调了 12 个参数,最后卡在 binlog 和 checkpoint

🔑 关键词:Flink CDC,MySQL,Iceberg,数据同步,调优

📖 摘要:基于真实环境踩坑:MySQL 8.0.34 单表 4.2 亿行同步到 Iceberg,Flink 1.18.1 全量+增量一开始跑 19 小时。文章拆解 binlog server-id 冲突、checkpoint 过密导致小文件、Kafka 分区不足等问题,给出 12 个参数和排查顺序,并对比全量 CDC、DataX+CDC、Canal 三条路线。

Flink CDC 同步 MySQL 到 Iceberg 太慢?我调了 12 个参数,最后卡在 binlog 和 checkpoint

图片

先交代环境,不交代环境谈调优就是耍流氓。MySQL 8.0.34,一主两从,binlog_format=ROW,binlog_row_image=FULL,binlog_transaction_dependency_tracking=WRITESET。单表 order_main 4.2 亿行,78 个字段,日增 1800 万。Flink 1.18.1 on YARN,CDC 2.4.2,Kafka 3.6.0 三节点,Iceberg 1.4.3 用 HadoopCatalog,HDFS 3.3.6。作业并行度一开始 8,TM 4 个,每个 16GB,RocksDB 增量 checkpoint,状态跑到 230GB。同步链路:MySQL -> Flink CDC -> Kafka -> Flink SQL -> Iceberg。目标 5 分钟内延迟,结果全量阶段跑了 19 小时还没完。

图片

先别急着加并行度。我先看 Flink UI 的 BackPressure,Source 端 high,下游 Iceberg sink 也 high。Metrics 里 numRecordsInPerSecond 只有 1.2 万,但 MySQL 侧 SHOW PROCESSLIST 看到 CDC 连接一直 Sending data。Kafka topic ods_mysql_binlog 当时 8 分区,副本 2,生产者吞吐 12MB/s,消费者 lag 最高 340 万。Checkpoint 间隔 10s,每次 18-25s,频繁超时。我犯的错:把 debezium.snapshot.fetch.size 设成 102400,以为越大越快,结果 MySQL 内存和网络被打满,其他业务报警。后来降到 10240,scan.incremental.snapshot.chunk.size=8096chunk-meta-cache-size=1024,才稳住。

图片

真正卡住的是两个地方。一个是 MySQL binlog,主库 max_binlog_size=1Gexpire_logs_days=7,但 Flink CDC 的 server-id 和 Canal 撞了,导致重复拉取和断连。我改成 5400-5404 这个范围,debezium.heartbeat.interval.ms=30000debezium.snapshot.locking.mode=none,避免了全量时锁表。另一个是 Iceberg 小文件。Checkpoint 10s 一次,每个 checkpoint 都提交一次 Iceberg,sink 并行度 8,结果每 10 秒生成 8 个 parquet 文件,一天 69120 个文件,NameNode 差点被干翻。把 execution.checkpointing.interval 从 10s 调到 60000,execution.checkpointing.min-pause=30sexecution.checkpointing.timeout=10minwrite.target-file-size-bytes=134217728write.distribution-mode=hashwrite.upsert.enabled=true,小文件降到每天 800 左右。

图片

参数不是越多越好。我列一下最后有效的 12 个:debezium.snapshot.fetch.size=10240scan.incremental.snapshot.chunk.size=8096debezium.snapshot.locking.mode=nonedebezium.heartbeat.interval.ms=30000server-id=5400-5404table.exec.source.idle-timeout=30sexecution.checkpointing.interval=60000execution.checkpointing.min-pause=30000execution.checkpointing.timeout=600000state.backend.incremental=truewrite.target-file-size-bytes=134217728sink.parallelism=16。Kafka 分区从 8 加到 16,副本 2,min.insync.replicas=1,消费 lag 从 340 万降到 3 万以内。但注意:sink.parallelism=16 后 TM 从 4 个加到 8 个,每个 8GB,总内存没变,反而 CPU 更稳。别盲目 32,Iceberg 写 HDFS 会抢带宽。

图片

方案对比也得说清楚。全量+增量一体化 Flink CDC,适合表少、单表小于 1000 万、能接受全量时主库压力;DataX 全量 + Flink CDC 增量,适合 500 张表以上、单表过亿,初始化走从库,增量再追;Canal + Kafka + Spark/Flink,适合已有 Canal 运维体系,但组件多。我的观点:如果业务能接受 10-30 分钟延迟,不要迷信流批一体,DataX 初始化 + CDC 增量更省心。我们最后把 1200 张表拆成:小表 800 张走 Flink CDC 全量+增量,大表 400 张走 DataX 从库全量 + Flink CDC 增量,整体延迟 3-8 分钟,主库 CPU 从 75% 降到 42%。

图片

最后说排查顺序。第一步看 Flink UI 反压和 checkpoint,第二步看 MySQL SHOW PROCESSLISTSHOW MASTER STATUS,第三步看 Kafka lag 和分区,第四步看 Iceberg 元数据和小文件,第五步才调并行度和内存。别一上来就加并行度,我试过从 8 加到 32,结果 checkpoint 更大,HDFS 写入冲突,作业反而重启 6 次。大数据开发不是写 SQL 就完事,是算资源、算失败路径、算恢复时间。你要是也卡在 CDC 同步,先把 checkpoint 间隔和 binlog server-id 查一遍,能省两天。