
早两年我做数据同步的时候方案还停留在“定时任务拉全量 业务代码双写 时间戳轮询”这种土办法上每次一到凌晨跑批就提心吊胆生怕任务超时或者数据对不上。后来在项目里引入了 Flink CDC才真正体会到什么叫“基于数据库日志做变更数据捕获”binlog 里每一条 insert、update、delete 都能被实时捕获、实时消费再配合 Flink 的分布式计算和状态管理整个数据同步这件事从“脚本”变成了一套工程化体系。这篇文章我想站在实际使用者的角度把 Flink CDC 的原理、选型、落地部署和排障经验完整讲一遍。不管你是准备做实时数仓还是想把业务库实时同步到 Kafka、搜索或分析型数据库都能从这里找到可直接参考的内容。1. 为什么同步方案最终选 Flink CDC1.1 传统同步方式的痛点我基本都踩过先聊聊那些年在同步上吃过的亏。第一种是定时全量同步每天凌晨用调度任务把业务库的表拉一遍到数仓。这个方案简单粗暴但问题也最明显延迟是按天算的业务上要求看实时数据时完全没法满足数据量一上来跑批窗口根本不够用经常是早上上班了任务还没跑完。第二种是业务系统双写在应用代码里同时写业务库和同步目标这个思路的问题在于侵入性太强而且一旦有一边写入失败两边数据就开始打架后面排查脏数据能排查到怀疑人生。还有一种常见做法是用时间戳轮询也就是根据表里的 update_time 字段定期查询最近变更的数据。这种方式看着聪明实际坑更多update_time 不是所有表都有即使有也可能被业务代码绕过delete 操作基本发现不了除非你额外做软删除标记而且每次轮询都是对业务库的一次全表扫描或者大范围索引查询数据库压力很大。触发器方案更不用说了在源库加触发器直接影响写入性能数据库 DBA 一般第一个反对。这些方法有一个共同点——它们都在“应用层”或者“数据本身”层面想办法离数据库内核太远了。而 CDCChange Data Capture的核心思路是直接从数据库的日志文件里读变更业务库什么都不用改对源库的影响也小得多。1.2 Flink CDC 在同类工具里的位置市面上做 CDC 的工具不少我接触过的就有 Canal、Debezium、Maxwell 这些。Canal 是 MySQL binlog 解析的成熟方案但主要围绕 MySQL 生态而且它本身不提供分布式计算和状态管理下游的断点续传、数据加工基本要自己再写一套。Debezium 本身是 JVM 框架对各种数据库支持都不错但它更像是为了 Kafka Connect 设计的要玩转它你得先熟悉 Kafka Connect 那一整套配置和运维体系。Maxwell 也是解析 MySQL binlog 的工具输出格式相对简单但分布式扩展和故障恢复能力都偏弱。Flink CDC 在这个赛道里最大的不同是它把 CDC 连接器直接打造成了 Flink 的 Source这意味着任务天然具备 Flink 的分布式并行能力、checkpoint 状态恢复能力和 SQL 表达能力。我做个简单的对比表你就明白了。方案源数据库支持分布式能力断点恢复下游生态Flink CDC主流的 MySQL、PostgreSQL、Oracle、SQLServer、MongoDB 等强基于 Flink 并行checkpoint 自动恢复Flink SQL / Kafka / 各类 SinkCanal主要 MySQL一般需自建要自己管理位点需要额外开发对接Debezium多数据库偏 JVM/Kafka Connect要结合外部框架依赖 Kafka Connect offset与 Flink 集成需额外封装Maxwell主要 MySQL有限要处理位点记录主要输出到 KafkaJSON 格式1.3 Flink CDC 最对口的几个场景结合我实际做过的项目Flink CDC 最适合这几类场景。第一类是实时数仓的 ODS 层建设业务库里几十张核心表通过 Flink CDC 实时同步到 Kafka后续 Flink SQL 从 Kafka 里消费做实时 ETL这条链路现在很多互联网公司都在用替代了原来跑批的功能。第二类是缓存和搜索更新比如订单状态一变需要立刻更新 Redis 缓存或者 Elasticsearch 索引用 Flink CDC 会让业务系统完全无感知不用在业务代码里到处埋点。第三类是数据库迁移和同步比如从自建 MySQL 迁移到云上数据库或者做跨机房的数据同步Flink CDC 可以先做全量快照再自动衔接增量 binlog整个过程业务几乎不用停机。还有一个比较实用的是微服务的事件驱动架构核心业务表的变化直接变成事件流其他服务消费事件做响应。2. 核心原理拆解从 binlog 到 Flink 无界流2.1 增量快照机制为什么可以做到基本不锁表Flink CDC 在早期全量阶段设计方案其实有妥协的地方就是和传统工具一样需要短暂地给表加锁来获取一致性快照点。传统做法是“FLUSH TABLES WITH READ LOCK”把整张表锁住拿到 binlog 位点然后开始导数据表越大锁的时间越长大库根本扛不住。Flink CDC 之所以体验好是因为它实现了增量快照Incremental Snapshot机制。它的思路是把一张表的主键范围拆成多个 chunk比如主键从 1 到 1 亿的表按每 4096 个主键区间拆成一个 chunk。每个 chunk 的读取都在一个短事务里完成事务只覆盖这一个区间拿到该区间数据的同时记录下当时的 binlog 位点然后立刻释放。这样多个 chunk 可以并行读取而且每个短事务的持有时间只有几秒甚至更短。所有 chunk 读完之后任务以这些 chunk 记录的最小 binlog 位点作为增量阶段的起点继续消费 binlog 里新产生的变更。由于全量阶段业务库的写入一直都在进行所以从最小位点开始重放 binlog 时部分数据可能已经包含在 chunk 快照里了会产生重复事件。但这保证了“不丢数据”重复问题依赖 Flink 的幂等机制或者下游系统去重。我在生产环境用下来的感受是这个机制真的很实用尤其面对几十亿行的大表时锁表时间基本可以忽略不计。2.2 一致性语义checkpoint 与 binlog 位点的配合如果一个同步任务跑着跑着挂了恢复之后怎么保证不重不漏靠的是 Flink 的 checkpoint 机制。Flink CDC Source 在做全量读取时会为每个 split也就是 chunk 或者 binlog 读取段记录当前的进度进入增量阶段后会周期性地把当前消费到的 binlog 位点保存到 state backend 里。当任务发生故障重启Flink 会自动从最近一次成功的 checkpoint 恢复所有算子状态。Source 会从保存的 binlog 位点继续读取而不是从头开始。这个能力非常重要因为我们实际运维时总会有各种意外网络抖动、集群资源不足导致任务重启、或者下游 Kafka 短暂不可用。如果没有 checkpoint 恢复机制任何一次意外都意味着手动记位点、手动补数想想都头疼。需要强调的是Flink CDC 提供的是 Source 侧的一致性保证也就是不会丢数据但由于全量阶段和增量阶段的衔接可能产生重复记录端到端的Exactly-Once语义还需要下游 Sink 配合。如果是写到 Kafka重复消息在大多数场景下可以接受下游消费时做幂等处理就行如果想把端到端一致性做得更彻底就需要配合 Flink 的两阶段提交 Sink。2.3 两代架构的价值分野2.x 的 SQL 连接器与 3.x 的 YAML PipelineFlink CDC 项目的发展可以明显分成两个阶段。2.x 时代最常用的方式是写 Flink SQL。你要先建一张映射业务表的 Source 表再建一张映射目标系统的 Sink 表最后写一条 INSERT INTO 语句把数据串起来。这个方式灵活适合和 Flink SQL 的实时计算逻辑混在一起用比如同步的同时做字段过滤、纬度关联、窗口聚合。3.x 版本推出了 YAML Pipeline 模式这是我最近用得越来越多的方式。它的基本思路是只写一个 YAML 文件声明 source 是什么、sink 是什么启动一个独立的同步任务。这个模式更像是专门给“数据同步”场景设计的你不用关心建表语句也不用拼 SQL上手门槛大大降低。而且 3.x 对表结构变更DDL 事件的支持更好也加入了整库同步能力有了更现代化的 schema 处理机制。你可以理解成2.x 让你拿 Flink SQL 玩出花3.x 让你一条命令完成整库同步。3. 从零搭建一个 MySQL 到 Kafka 的实时同步任务3.1 环境准备和版本选型在真正写配置之前先把环境准备扎实。你需要 JDK 1.8 或 11、一个可用的 Flink 运行环境本地 Standalone 就行集群也可以、MySQL 开启 binlog、以及 Kafka 2.x 以上版本。MySQL 的 binlog 必须改成 ROW 模式只记录变更行的前后镜像否则解析不出字段级变化。还需要确认 binlog_row_image 设置为 FULL这样 binlog 里才能拿到完整的旧值和新值方便下游处理。权限这块通常容易被忽略。Flink CDC 账号除了要能 SELECT 表数据还必须拥有 REPLICATION SLAVE 和 REPLICATION CLIENT 权限这样才能通过 MySQL 复制协议读取 binlog 元数据和位点。我一般创建账号时的授权语句长这样CREATE USER cdc_user% IDENTIFIED BY your_password; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc_user%; FLUSH PRIVILEGES;版本选型要提醒一句不要盲目上最新版本。不同 Flink CDC 版本对 Flink 版本有明确的适配范围用之前一定要去官方文档查兼容性矩阵。我生产环境比较常用的是 2.4.x 搭配 Flink 1.17整体稳定主流连接器也都覆盖了。如果你想用 3.x 的 YAML Pipeline 模式建议先在小范围业务验证确认和你的 Flink 集群版本匹配后再推全量。3.2 方式一用 Flink SQL 完成单表同步用 Flink SQL 做单表同步本质是定义三件事源表、目标表、同步逻辑。我这里用一个订单表同步到 Kafka 的例子先创建 Source 表CREATE TABLE orders_source ( id INT, user_id BIGINT, product_id BIGINT, amount DECIMAL(10, 2), status VARCHAR(20), create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 127.0.0.1, port 3306, username cdc_user, password your_password, database-name shop, table-name orders, server-id 5401, scan.incremental.snapshot.chunk.size 4096 );再创建 Sink 表这里用 Kafka 的 debezium-json 格式这样消息体里会包含 before 和 after 的完整结构便于下游消费时感知变更前后状态CREATE TABLE orders_kafka ( id INT, user_id BIGINT, product_id BIGINT, amount DECIMAL(10, 2), status VARCHAR(20), create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector kafka, topic ods_orders, properties.bootstrap.servers 127.0.0.1:9092, properties.group.id order_sync_group, format debezium-json );最后执行同步语句INSERT INTO orders_kafka SELECT id, user_id, product_id, amount, status, create_time, update_time FROM orders_source;这里有两个参数值得展开。server-id 必须保证在整个 MySQL 主从复制拓扑里唯一否则会被 MySQL 判定为重复 slave 连接直接踢掉。scan.incremental.snapshot.chunk.size 决定了全量阶段每个 chunk 的行数不是越大越好我实测下来 4096 左右比较均衡兼顾了读取效率和源库压力。3.3 方式二用 3.x YAML Pipeline 完成整库同步如果你不关心 SQL 里的复杂转换只想快速把整个库或者一批表同步到 KafkaYAML Pipeline 模式很合适。你只需要写一个类似这样的配置文件source: type: mysql hostname: 127.0.0.1 port: 3306 username: cdc_user password: your_password tables: shop.orders, shop.products server-id: 5401-5404 sink: type: kafka properties.bootstrap.servers: 127.0.0.1:9092 topic: ods_orders format: debezium-json pipeline: name: MySQL to Kafka Sync parallelism: 2启动命令也很简单bin/flink-cdc.sh pipeline /path/to/orders.yaml这个模式最大的好处是不需要写建表语句不需要拼 INSERT SQL连接器会自动读取表结构并维护元数据。server-id 写成 5401-5404 这种范围是因为 3.x 的并行 task 会从范围里拿不同的 server-id 去连接 MySQL避免互相冲突。如果你的表数量比较多甚至可以用 shop.* 这种通配方式做整库同步一旦新增表任务能自动感知并开始同步这在传统方式里是很难实现的。3.4 数据链路验证从新增到删除逐一验证配置完成后最怕的是任务显示运行中实际数据没进去。我习惯按下面这个顺序验证全链路。先用 Kafka 客户端工具看 topic 是否存在kafka-topics.sh --bootstrap-server 127.0.0.1:9092 --list | grep ods_orders然后启动一个消费者准备观察消息kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 --topic ods_orders --from-beginning接着在 MySQL 里依次执行三种操作。先插入一条数据观察 Kafka 里是否出现 after 字段有值的 JSON 消息再执行一次 UPDATE观察消息里 after 变化而 before 保留旧值最后执行 DELETE确认消息里 before 有值、after 为 null。如果这三类消息都能正常出现说明 binlog 解析、格式转换、Kafka 写入这条链路是完全通的。4. 高频问题与排障心得4.1 我碰过的坑整理成一张排查表实际使用 Flink CDC 这一年多我把碰过的问题按现象、原因和处理方式梳理成了一个速查表遇到问题时能少走很多弯路。现象大概率原因处理方式任务启动报 Access denied账号缺少 REPLICATION 相关权限补授权并确认 binlog 开启连接 MySQL 报 server-id 冲突多个任务共用了同一 server-id每个任务分配独立 server-id 或范围无主键表全量阶段卡住无法拆 chunk回退到锁表模式给表加主键或指定 chunk key-column时间字段和源库相差 8 小时时区设置不一致Debezium 默认输出相关时区统一 Flink 和连接器时区配置表结构变更后任务异常SQL 任务无法自动感知 schema 变化更新建表语句并重启任务或用 3.x Pipeline同步延迟持续变大源库大事务、并行度过低或下游写入慢检查大事务、提高并行度、查看 checkpoint 耗时4.2 server-id 冲突是最容易踩的雷server-id 问题我至少见过三次了而且每次发生时的表现都不太一样。有时候是任务反复断线重连日志里出现 MySQL 主从复制的报错有时候是任务表面正常但数据同步突然停止更新过段时间才恢复。根本原因都是同一个多个同步任务或者同一个任务里的多个并发线程使用了相同的 server-id去连接 MySQL binlog。MySQL 主从复制协议里面每个 slave 在 master 看来都有一个唯一标识。Flink CDC 本身扮演的其实就是一个“伪 slave”角色如果它的 server-id 和别的从库或者其他 CDC 任务一样master 只能保留最后那个连接前面的连接就会被强制断开。正确做法是给每一个 Flink CDC 任务分配一个独立的 server-id需要并行读取多个表时分配一段连续范围。比如你任务并行度是 2可以写 server-id5401-5402如果并行度是 4就写 5401-5404。范围越大意味着同一时间打开的复制连接越多所以也别滥用。4.3 大表同步的性能调优经验大表场景下全量快照往往会跑几个小时这时候有几个参数需要配套调优。scan.incremental.snapshot.chunk.size 决定 chunk 拆分的粒度表特别大时可以适当调大比如 8096减少 chunk 数量从而降低调度开销但如果调得太大单个 chunk 的短事务执行时间变长对源库压力也变大。建议先根据表行数和行大小估算一下比如 1 亿行的表400 万行一个 chunk大概会产生 25 个任务切片。另外一个很重要的点是 binlog 保留时间和全量快照耗时的关系。因为全量快照读取期间业务库的写入并没有停binlog 在持续产生所有 chunk 读完以后增量阶段从最小位点开始消费如果此时最早的 binlog 已经被清理任务会因为找不到位点而报错甚至无法启动。所以在启动大表同步之前一定要确认 binlog 保留时间大于预估全量耗时。我一般会先跑一个测试任务观察全量阶段大概需要多长时间然后再确认 binlog 保留策略是否覆盖得住。4.4 表结构变更的应对策略表结构变更DDL是 CDC 同步里绕不开的话题。在 2.x 的 SQL 模式下如果业务表新增了一列同步任务不会自动感知甚至可能因为 schema 校验不匹配导致反序列化失败任务卡住。处理方式比较传统先在下游目标系统里手动加上对应字段然后修改 Flink SQL 里的建表语句最后重启任务并从最近 checkpoint 恢复。3.x 的 YAML Pipeline 模式在这方面有改进它会把源端的 DDL 事件也转成事件消息发送到下游比如 Kafka 里会出现一条 schema change 事件下游消费端可以据此自动加列或调整结构。这在整库同步场景很实用。不过要记住schema 自动演进依赖下游系统的接收能力不是所有目标都支持用之前要先确认好。我的建议是业务有大规模 DDL 计划时还是提前和开发对一遍变更内容把同步任务调整放在维护窗口里一起做不要在业务高峰硬扛。5. 典型落地场景实时数仓里 Flink CDC 的位置5.1 一条链路讲清楚 ODS 层的同步逻辑很多团队的实时数仓架构里Flink CDC 处于最上游、最基础的位置。业务系统的订单、用户、支付流水这些数据都在 MySQL 里通过 Flink CDC 同步到 Kafka 对应的 ODS 主题。Kafka 之后接 Flink SQL 做实时清洗、订单维度关联、累积窗口统计最后写入 Doris 或者 ClickHouse 这样的分析型数据库供报表查询。这个过程的好处在于业务系统完全不需要改动同步任务独立部署新增一张表只要在 CDC 任务里加一个配置几分钟后数据就能实时出现在数仓里。相比过去每天凌晨跑批实时数仓的查询结果能达到秒级刷新业务方对数据时效性的满意度提升明显。你可以把 Flink CDC 理解为水流入口的总闸闸门打开后数据源源不断流进管道后面每级处理单元的延迟和吞吐都可以单独调节。我在实际项目中还见过一种组合玩法一个 Flink CDC 任务同时同步多个库的表在 Kafka 里按业务域拆分成不同 topic后面的实时任务再按需订阅。这样既节省了源库连接资源又保持了逻辑上的清晰。5.2 同步任务上线前我建议做这三件事第一件事是确认源库的 binlog 保留时长和磁盘空间。很多生产环境的 binlog 只保留 24 小时如果你的同步任务因为 Kafka 挂了或者其他原因停止了一天以上恢复后根本找不到开始的 binlog 位点只能重新全量同步。我会在评估阶段就找运维确认数据保留策略能调长就调长几天给自己留出容错空间。第二件事是压测一次大事务。不要只在测试环境验证简单的增删改查真实业务里很可能出现一个 UPDATE 语句影响几十万行的情况binlog 会生成一个巨大的事务体Flink CDC 消费这种大事务时常会有一段明显延迟。我在上线前会专门用一条数据更新脚本对表做大更新观察同步任务的延迟曲线和内存占用确保没有明显瓶颈。第三件事是设置合理的 checkpoint 策略。Flink CDC 任务的生产配置里checkpoint 间隔不能太短否则频繁做状态快照会增加负担也不能太长否则任务重启后恢复时间偏长也可能丢更多未保存的进度。我一般设置 30 秒到 60 秒状态后端使用 RocksDB保证状态量大的时候也能稳定运行。5.3 用下来的真实体会我在实际使用中最受益的是把 Flink CDC 当成一个真正的“流系统”去设计而不是简单地把它当同步工具用。它的 Source 是流Sink 是流中间还能做各种处理整个任务的监控指标也都暴露在 Flink 体系里。每次看着 JobManager 日志里 binlog 位点持续向前推进、checkpoint 稳定完成那种踏实感是以前凌晨跑批完全比不了的。最后分享一个小技巧给每个同步任务命名时把源库、目标、业务模块都写清楚比如 mysql_shop_orders_to_kafka。任务多了以后你会感谢这个好习惯。真正把一套 CDC 链路跑稳靠的不是某个参数调得多完美而是每一步都留有冗余、每一类异常都有预案。希望这些经验和踩坑记录能让你在接手 Flink CDC 时少走几段弯路。