Flink CDC 3.x 实战:构建 MySQL 到 Iceberg 的实时数据同步管道

发布时间:2026/8/23 5:57:24
Flink CDC 3.x 实战:构建 MySQL 到 Iceberg 的实时数据同步管道 这次我们来看一个 Flink CDC 3.x 的技术项目。如果你正在处理数据库实时同步、构建实时数仓或数据湖并且对如何高效、低延迟地捕获和处理数据库变更数据感兴趣那么这篇文章值得你花时间读完。Flink CDC 3.x 是 Apache Flink 社区推出的新一代 Change Data Capture 框架它核心解决的是将数据库的变更事件如 MySQL 的 INSERT、UPDATE、DELETE实时、准确地接入到 Flink 流处理引擎中并最终流向数据湖仓如 Iceberg、Hudi或其他下游系统。最值得关注的几个特点是连接器原生集成无需额外部署 Debezium 等组件、全增量一体化读取自动切换快照和增量阶段、更低的源端压力通过并行读取和无锁算法优化以及对实时湖仓一体场景的深度支持。对于开发者而言这意味着更简单的部署、更稳定的同步链路和更强的实时数据处理能力。本文将带你快速了解 Flink CDC 3.x 的核心能力并完成从环境准备、连接器部署、到数据同步与入湖的完整实操验证。你会看到如何配置一个从 MySQL 到 Apache Iceberg 的实时同步管道并观察其运行状态和资源占用。无论你是想评估技术选型还是准备在生产环境落地这篇文章提供的步骤和排查思路都能作为参考。1. 核心能力速览在深入细节之前我们先通过一个表格快速把握 Flink CDC 3.x 的关键信息这有助于你判断它是否适合你的场景。能力项说明项目类型Apache Flink 的 Change Data Capture 连接器框架核心功能实时捕获数据库变更CDC支持全量增量一体化同步数据入湖/入仓推荐运行模式Flink SQL / Table API集成在 Flink 运行时中主要支持数据库MySQL, PostgreSQL, MongoDB, Oracle, SQL Server 等目标系统支持Apache Iceberg, Apache Hudi, Apache Paimon, Kafka, 各类数据库等部署复杂度中等。需准备 Flink 集群和对应的 Connector JAR。是否支持 API主要通过 Flink SQL/Table API 进行配置和提交作业提供丰富的 Connector 参数。是否支持批量任务是。全量读取阶段本质是批量快照增量阶段为持续流式任务。支持整库同步等“批量”表同步场景。资源占用关键点取决于数据源表数量、数据量、QPS。并行度、Checkpoint 间隔、Lakehouse 表格式的合并小文件操作会影响 CPU、内存和 IO。适合场景数据库实时同步至数据湖/仓、实时数仓 ETL、微服务数据聚合、归档与查询分离。2. 适用场景与使用边界Flink CDC 3.x 并非万能工具明确其适用边界能帮助你做出更合适的技术决策。它非常适合以下场景实时数仓/湖仓增量构建需要将业务库如 MySQL、PostgreSQL的数据变更实时同步到 Iceberg、Hudi 等湖仓表中为下游即席查询或批处理提供新鲜数据。微服务数据聚合多个微服务数据库的数据需要实时汇聚到一个中心化的分析库中进行关联分析和统一查询。缓存与搜索索引更新数据库变更需要实时驱动 Redis 缓存更新或 Elasticsearch 索引重建。审计与合规需要实时捕获所有数据变更记录用于审计追踪或满足数据合规性要求。传统 ETL 链路简化替代基于定时批量查询的 ETL 作业实现更低延迟、更少资源消耗的增量数据集成。它可能不是最佳选择或需要注意的边界极简一次性数据迁移如果只需要做一次性的全量数据迁移使用传统的mysqldump或 DataX 等工具可能更简单直接。源端数据库版本与权限Flink CDC 连接器对数据库版本有要求如 MySQL 需要 5.7 并开启 binlog且需要具有读取 binlog 或类似日志的权限。生产环境需与 DBA 协作。数据结构频繁变更如果源表频繁进行 DDL 操作如增删改列虽然 Flink CDC 支持 Schema Evolution但下游湖仓表的 Schema 变更处理需要额外配置和测试可能带来复杂度。无序或延迟容忍度极低CDC 数据在极端情况下可能存在短暂延迟或由于网络分区导致乱序。对于要求强一致性和毫秒级延迟的金融交易场景需要结合事务日志等更复杂的方案。运维与监控成本引入一个实时流处理框架意味着需要维护 Flink 集群的稳定性监控作业状态、延迟、反压等指标这需要一定的运维能力。合规与安全提醒使用 CDC 技术同步生产数据时务必确保符合数据安全法规。同步敏感数据如用户个人信息时应考虑在链路中集成数据脱敏组件。同时确保拥有操作源数据库的合法授权并限制下游数据存储的访问权限。3. 环境准备与前置条件在开始部署和测试之前请确保你的环境满足以下基本要求。我们将以一个典型的MySQL - Apache Iceberg的同步场景为例。操作系统Linux (CentOS 7, Ubuntu 18.04) 或 macOS。Windows 可用于开发测试生产环境建议 Linux。JavaApache Flink 需要 Java 8 或 Java 11。确保JAVA_HOME环境变量已正确设置。java -versionApache Flink需要准备 Flink 运行时。可以从 Apache Flink 官网 下载。本文示例基于Flink 1.17或1.18版本。你可以选择Standalone 集群用于快速测试。Flink on YARN/K8s用于生产环境。数据库 (以 MySQL 为例)版本 5.7 或 8.0。必须开启 Binlog并且格式为ROW。授予 Flink 连接用户足够的权限SELECT,RELOAD,SHOW DATABASES,REPLICATION SLAVE,REPLICATION CLIENT。配置server-id和binlog相关参数。一个简化的my.cnf配置示例如下[mysqld] server-id 1 log_bin /var/log/mysql/mysql-bin.log binlog_format ROW expire_logs_days 10 max_binlog_size 100M binlog_row_image FULL数据湖格式 (以 Apache Iceberg 为例)需要准备一个兼容 Iceberg 的 Catalog 后端如Hadoop Catalog依赖 HDFS或JDBC Catalog依赖关系型数据库如 MySQL/PostgreSQL 存储元数据。本文使用 Hadoop Catalog 进行演示因此需要Hadoop 环境或至少能访问 HDFS/对象存储和Iceberg Flink Runtime JAR。网络运行 Flink 任务的机器需要能访问源端 MySQL 和目标端存储如 HDFS。4. 安装部署与启动方式Flink CDC 3.x 的核心是 Connector JAR 包。部署的本质是将这些 JAR 包放入 Flink 的lib目录然后通过 SQL Client 或编程方式提交作业。步骤 1下载 Flink 和 Connector下载 Apache Flink 发行包并解压。wget https://dlcdn.apache.org/flink/flink-1.17.2/flink-1.17.2-bin-scala_2.12.tgz tar -xzf flink-1.17.2-bin-scala_2.12.tgz cd flink-1.17.2下载 Flink CDC 3.x 连接器。以flink-sql-connector-mysql-cdc-3.0.0.jar为例可以从 Maven 中央仓库或 Flink 官方仓库获取放置到lib/目录。wget -P lib/ https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/3.0.0/flink-sql-connector-mysql-cdc-3.0.0.jar下载 Iceberg Flink Runtime JAR版本需与 Flink 和 Iceberg 兼容同样放入lib/目录。wget -P lib/ https://repo1.maven.org/maven2/org/apache/iceberg/iceberg-flink-runtime-1.17/1.3.0/iceberg-flink-runtime-1.17-1.3.0.jar步骤 2启动 Flink 集群进入 Flink 目录启动一个本地 Standalone 集群。# 启动集群 ./bin/start-cluster.sh # 检查是否启动成功访问 Web UI (默认 http://localhost:8081) # 也可以通过jps查看进程 jps你应该能看到StandaloneSessionClusterEntrypoint和TaskManagerRunner进程。步骤 3启动 SQL ClientFlink 提供了 SQL Client 工具方便我们通过 SQL 与集群交互。./bin/sql-client.sh启动后你会进入 SQL 命令行界面。5. 功能测试与效果验证现在我们开始构建一个从 MySQL 到 Iceberg 的完整实时同步链路。整个过程分为创建源表映射 MySQL、创建目标表映射 Iceberg、执行同步作业。5.1 准备测试数据源MySQL在 MySQL 中创建一个测试数据库和表并插入一些数据。-- 在 MySQL 中执行 CREATE DATABASE test_cdc; USE test_cdc; CREATE TABLE user_actions ( id INT PRIMARY KEY, user_id INT, action VARCHAR(50), action_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ); INSERT INTO user_actions (id, user_id, action) VALUES (1, 1001, login), (2, 1002, purchase), (3, 1001, logout);5.2 在 Flink SQL 中定义 MySQL CDC 源表在 Flink SQL Client 中执行以下 DDL 来创建一个虚拟表它指向我们刚创建的 MySQL 表。-- 在 Flink SQL Client 中执行 CREATE TABLE mysql_user_actions ( id INT, user_id INT, action STRING, action_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname your_mysql_host, -- 替换为你的 MySQL 地址 port 3306, username your_username, -- 替换为你的用户名 password your_password, -- 替换为你的密码 database-name test_cdc, table-name user_actions, server-time-zone Asia/Shanghai );关键参数说明connector mysql-cdc指定使用 MySQL CDC 连接器。scan.startup.mode默认为initial表示先做全量快照再持续读取增量。这是“全增量一体”的关键。server-time-zone建议设置避免时区问题。5.3 在 Flink SQL 中定义 Iceberg 目标表接下来定义一个 Iceberg 表来接收数据。这里假设你的 Hadoop 环境或 MinIO 等 S3 兼容存储已就绪并配置了warehouse路径。-- 在 Flink SQL Client 中执行 -- 首先创建一个 Iceberg Catalog CREATE CATALOG iceberg_catalog WITH ( typeiceberg, catalog-typehadoop, warehousehdfs://localhost:9000/user/iceberg/warehouse -- 替换为你的 warehouse 路径 ); -- 使用该 Catalog USE CATALOG iceberg_catalog; -- 创建数据库和表 CREATE DATABASE IF NOT EXISTS test_db; USE test_db; CREATE TABLE if not exists iceberg_user_actions ( id INT, user_id INT, action STRING, action_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( format-version 2, write.upsert.enabled true -- 启用 Upsert根据主键合并变更 );关键参数说明catalog-typehadoop使用 Hadoop Catalog 管理 Iceberg 元数据。warehouseIceberg 数据文件和元数据的存储根路径。write.upsert.enabled true这对于 CDC 数据至关重要。它确保根据主键id对UPDATE和DELETE操作进行合并最终在 Iceberg 表中呈现最新的数据快照。5.4 启动实时同步作业现在通过一个INSERT INTO语句将源表的数据持续同步到目标表。这个语句本身就是一个永不停歇的 Flink 流作业。-- 在 Flink SQL Client 中执行 INSERT INTO iceberg_catalog.test_db.iceberg_user_actions SELECT * FROM default_catalog.default_database.mysql_user_actions;提交后Flink 作业就开始运行了。你可以通过 Flink Web UI (http://localhost:8081) 看到这个运行中的作业。5.5 验证同步效果初始全量同步作业启动后会立即读取user_actions表的当前全量数据3条记录并写入 Iceberg。你可以通过查询 Iceberg 表来验证。SELECT * FROM iceberg_catalog.test_db.iceberg_user_actions;增量变更捕获回到 MySQL执行一些增、删、改操作。-- 在 MySQL 中执行 INSERT INTO user_actions (id, user_id, action) VALUES (4, 1003, view); UPDATE user_actions SET action buy WHERE id 2; DELETE FROM user_actions WHERE id 3;实时查询验证稍等片刻通常秒级再次查询 Iceberg 表。你应该能看到新增了id4的记录。id2的记录action字段从purchase更新为buy。id3的记录被标记为删除对于 Upsert 表查询最新快照时这条记录不会出现。 这证明了 Flink CDC 实时捕获并处理了 MySQL 的变更。6. 接口 API 与批量任务虽然我们主要通过 SQL 交互但 Flink CDC 的“批量”能力体现在作业的启动模式和任务管理上。全量读取作为“批量”阶段在scan.startup.mode initial模式下作业启动时会先执行一个全量快照读取。对于大表这个阶段类似于一个批处理任务Flink CDC 3.x 支持并行快照读取可以加速这个过程。你可以通过scan.incremental.snapshot.chunk.size等参数来优化全量阶段的性能。整库同步与表发现Flink CDC 3.x 支持整库同步即自动捕获整个数据库内所有表的变更。这本质上是由一个作业管理多个表的“批量”同步任务。配置时使用table-name .*进行正则匹配即可。这极大地简化了多表同步的运维成本。作业管理与 API生产环境中我们通常通过REST API或Flink CLI来管理作业提交、停止、保存点。例如你可以将上面的 SQL 语句写在一个.sql文件中然后通过 SQL Client 的-f参数或flink run命令提交。# 将 SQL 语句保存在 sync_job.sql 文件中然后提交 ./bin/sql-client.sh -f /path/to/sync_job.sql # 或者如果你将作业打包成 JAR可以使用 flink run ./bin/flink run -c com.example.StreamingJob /path/to/your-job.jar通过 Flink Web UI 或 REST API (http://localhost:8081/jobs)你可以监控所有作业的状态、背压、吞吐量等指标这对管理批量同步任务至关重要。7. 资源占用与性能观察运行 Flink CDC 作业时需要关注以下几个方面的资源使用情况TaskManager 内存这是 Flink 执行任务的内存。在conf/flink-conf.yaml中配置。对于 CDC 作业尤其是同步多张大表或高 QPS 表时需要分配足够的堆内存和堆外内存。建议从 2-4G 开始测试根据实际情况调整。taskmanager.memory.process.size: 4096m # TaskManager 总进程内存并行度与 CPU源表读取的并行度parallelism.source会影响 CPU 使用率。对于单表通常设置为 1对于整库同步或多表可以适当提高。通过 Web UI 的 “作业概览” 可以观察各个算子的繁忙程度。Checkpoint 与状态存储Flink CDC 需要定期做 Checkpoint 来保证 Exactly-Once 语义。Checkpoint 间隔execution.checkpointing.interval会影响状态大小和 IO。间隔太短会增加负载太长则恢复时间变长。状态后端如 RocksDB会占用本地磁盘空间。目标端写入压力写入 Iceberg/Hudi 时频繁提交会产生大量小文件。需要合理配置write.parquet.row-group-size-bytes、write.target-file-size-bytes以及流式写入的提交间隔平衡写入性能和查询性能。网络与连接数Flink CDC 源任务会与 MySQL 建立数据库连接来读取 Binlog。确保 MySQL 的max_connections配置足够并监控连接数。监控建议充分利用 Flink Web UI 监控反压指标、吞吐量numRecordsInPerSecond和 Checkpoint 时长。如果发现反压可能是下游写入如 Iceberg Commit或网络成为瓶颈需要针对性优化。8. 常见问题与排查方法在部署和运行 Flink CDC 时你可能会遇到以下问题。这里提供基本的排查思路。问题现象可能原因排查方式解决方案作业启动失败报错连接不上数据库1. 网络不通或防火墙限制。2. 数据库地址、端口、用户名、密码错误。3. MySQL 未开启 Binlog 或格式不对。4. 用户权限不足。1. 使用telnet或mysql客户端测试连接。2. 检查 DDL 中的 WITH 参数。3. 在 MySQL 中执行SHOW VARIABLES LIKE ‘%binlog%’;。4. 检查用户权限SHOW GRANTS FOR ‘user’‘host’;。1. 开通网络。2. 修正连接参数。3. 修改my.cnf并重启 MySQL。4. 授予所需权限。全量同步阶段卡住或无进度1. 表数据量极大快照慢。2. 源表有未提交的长事务或大锁。3. 并行度设置不合理。1. 观察 Web UI 中 Source 算子的记录数是否增长。2. 在 MySQL 中查询SHOW PROCESSLIST;和SELECT * FROM information_schema.innodb_trx;。3. 检查作业并行度。1. 耐心等待或调整chunk-size。2. 结束无关长事务。3. 对于大表可尝试增大源并行度需 MySQL 配置支持。增量同步延迟高1. 下游 Sink如 Iceberg Commit慢。2. 网络抖动。3. Flink 作业出现反压。1. 在 Web UI 查看 Sink 算子的背压状态红色表示高反压。2. 监控 Flink 的currentFetchEventTimeLag指标。3. 检查目标存储如 HDFS的 IO 状态。1. 优化 Iceberg 写入参数如增大文件大小调整提交间隔。2. 检查网络。3. 增加 Sink 并行度或调优 Checkpoint。Iceberg 表查询不到最新数据1. Flink 作业未正常运行。2. Iceberg 表格式版本或 Catalog 配置问题。3. 数据未成功提交。1. 检查 Flink 作业状态是否为RUNNING。2. 检查 Iceberg 表的元数据文件是否更新。3. 在 Flink 日志中搜索Committed snapshot相关日志。1. 重启或恢复作业。2. 确认使用的 Iceberg 版本与 Flink Runtime 兼容。3. 确保write.upsert.enabled等参数配置正确。作业失败后重启数据重复或丢失1. Checkpoint 未成功启用或配置错误。2. 使用了‘scan.startup.mode’ ‘latest-offset’跳过了失败前的 Binlog 位置。1. 检查flink-conf.yaml中状态后端和 Checkpoint 配置。2. 检查作业日志中的 Checkpoint 完成记录。3. 审查启动模式。1. 确保配置了有效的状态后端如 RocksDB和合理的 Checkpoint 间隔。2. 生产环境建议使用‘initial’模式并从保存点Savepoint恢复。9. 最佳实践与使用建议基于测试和常见问题这里总结一些提升稳定性和效率的建议测试先行在生产环境部署前务必在准生产环境进行全链路测试包括全量同步性能、增量同步延迟、源端 DDL 变更、作业重启恢复、网络中断模拟等。配置优化Checkpoint根据数据量和延迟要求设置合理的间隔如 1-5 分钟。太短会增加负载太长则恢复慢。并行度根据源表数量和下游处理能力设置。不建议盲目调高。MySQL 侧确保binlog_row_imageFULL为连接用户赋予足够权限并监控 Binlog 空间。监控告警建立关键指标监控如作业状态、消费延迟sourceIdleTime、反压情况、Checkpoint 失败率、MySQL 连接数等。集成到公司现有的监控告警体系如 Prometheus AlertManager。数据一致性保障启用Exactly-Once语义依赖于 Checkpoint 和下游 Sink 的支持。对于 Iceberg/Hudi 等支持 Upsert 的目标务必开启‘write.upsert.enabled’‘true’以正确处理更新和删除。定期对比源库和目标表的数据一致性尤其是在作业发生故障切换后。版本与兼容性管理记录并锁定 Flink、Flink CDC Connector、目标 Connector如 Iceberg、数据库驱动等各组件的版本避免因版本升级带来的不兼容问题。Schema 变更处理如果源表结构可能变更需提前规划 Schema Evolution 策略。测试 Flink CDC 的‘scan.newly-added-table.enabled’等参数并了解下游 Iceberg 表对新增列、修改列类型的支持情况。安全与权限生产环境中使用专属的、权限最小化的数据库账户。对存储于 HDFS/S3 的 Iceberg 数据文件进行访问权限控制。同步链路中如需传输敏感数据考虑在 Flink 作业中集成加密或脱敏 UDF。Flink CDC 3.x 将数据库实时同步的门槛显著降低其“全增量一体”和“无锁读取”的特性让实时数据集成变得更加优雅。对于构建实时湖仓来说它提供了一个稳定、高效的源头活水。建议你先在一个简单的测试环境单节点 Flink 测试 MySQL 本地文件系统模拟 Iceberg中走通整个流程感受从 Binlog 变化到湖仓表查询的端到端延迟。之后再根据业务的数据量、表数量和运维能力逐步将其应用到更复杂的生产场景中。