☰
Flink 1.17 + Flink-CDC 2.4 实现 MySQL CDC 实时同步 Demo 全解析
2026/9/26 6:05:33 网站建设 项目流程

提个问题:你在网上搜“Flink MySQL CDC Demo”,搜出来的十篇里至少八篇跑不起来。不是Flink版本老到连类都换了,就是CDC连接器的Maven坐标抄错,还有一半人折腾了半天最后卡在“连上了却没增量数据”——因为根本没人提MySQL端要开binlog和配权限。这套东西我整理过很多次,最后固定下来的组合很简单:Flink 1.17 + Flink-CDC 2.4 + Java,本地IDE直接跑main方法,不需要额外搭Flink集群,MySQL里插入一条记录,控制台立刻输出一条JSON。这篇文章就把这套“复制就能跑”的Demo完整拆给你,包含环境版本匹配逻辑、MySQL服务端配置、核心代码解读,以及我实测中踩过的坑位清单。

适合谁来读?刚接触CDC想快速跑通原理的,被网上版本混搭折磨到想放弃的,还有准备拿Flink CDC做实时数仓、但想先起一个最小闭环验证可行性的。内容按“为什么这样选型、代码怎么组织、数据库怎么配合、跑起来怎么看效果、出了问题怎么排查”的顺序展开,你可以直接照着抄。

1. 先把版本关系理清楚:Flink 1.17 配 Flink-CDC 2.4 为什么稳

1.1 版本坐标变化的坑:com.ververica 还是 org.apache.flink

Flink CDC的项目坐标在2.x和3.x之间发生过一次“搬家”。2.4及之前的版本,groupId是com.ververica,artifactId是flink-connector-mysql-cdc;到了3.0之后,整个项目进入Apache Flink组织,groupId变成了org.apache.flink,artifactId也改成了flink-cdc-connector-mysql。

这个变化坑了很多人:从老文章里复制了2.x的坐标,读的却是3.x的文档,代码里类名对不上,编译直接报错。反过来,用3.x的坐标去跑2.4的代码,也会出现工厂类找不到的问题。本文固定使用2.4.x,坐标一定认准com.ververica这条线,别混。

1.2 组件兼容矩阵与选型逻辑

我自己在多个版本组合里实测过,列个表给你参考:

组件推荐版本说明
Flink1.17.2官方稳定版,Java 8和Java 11都支持
Flink-CDC MySQL连接器2.4.2官方声明的兼容范围包含Flink 1.15~1.17
MySQL5.7或8.0必须开启binlog且格式为ROW
JDK8或11建议11,别用17去跑Flink 1.17
Maven3.6+工程构建工具,没什么特殊要求

这套组合最大的优势是“文档密度高”。Flink 1.17是社区使用最广的版本之一,Flink CDC 2.4的API设计又保留了最直观的DataSource写法,你在网上搜到的大部分报错和解决方案都能直接对应上,对新手非常友好。

1.3 为什么暂时不用Flink-CDC 3.0

Flink CDC 3.0把重心转向了Pipeline模式,用一份YAML定义整条同步链路,思路很先进。但对一个“跑通原理”的入门Demo来说,2.4的编程式API反而更合适——MySqlSource.builder()链式调用,参数一目了然,每条配置都能对应到MySQL端的一个真实概念。等你看懂了binlog位点、理解了Debezium的ChangeEvent结构,再去上手3.0的Pipeline模式会轻松很多。

2. 复制就能跑的核心代码:环境准备与完整工程

2.1 本地环境要求与两个容易绕弯的点

跑这套Demo不需要单独安装Flink,也不需要启动任何集群。Flink本身是一套Java库,StreamExecutionEnvironment.getExecutionEnvironment()在IDE里执行时会自动以Local模式运行,这一点比Spark省心很多。

环境上注意两点:

一是JDK版本。Flink 1.17官方支持Java 8和Java 11,但别用Java 17跑。原因很现实:Flink内部不少第三方依赖在Java 17下会触发模块化访问限制,报一些“InaccessibleObjectException”,排查起来很浪费时间。用Java 11最省心。

二是Maven的maven.compiler.source/target配置。如果你本机的JDK是17,而项目pom里写的是8,编译时会提示“源发行版 8 需要目标发行版 8”之类的错误。这个跟Flink无关,纯粹是Maven编译器插件和JDK版本不匹配,把pom里的source/target设成和JDK一致就行。我见过有人在热词里搜“java: 警告: 源发行版 17 需要目标发行版 17”,就是这类问题。

2.2 pom.xml 完整依赖清单

直接新建一个Maven工程,把下面的依赖贴进pom.xml:

<properties> <flink.version>1.17.2</flink.version> <flink.cdc.version>2.4.2</flink.cdc.version> <maven.compiler.source>11</maven.compiler.source> <maven.compiler.target>11</maven.compiler.target> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> </properties> <dependencies> <!-- Flink DataStream API --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>${flink.version}</version> </dependency> <!-- 本地运行所需,提供LocalEnvironment入口 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>${flink.version}</version> </dependency> <!-- 如果后续想用SQL方式建CDC表,需要table api的bridge --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-api-java-bridge</artifactId> <version>${flink.version}</version> </dependency> <!-- Flink CDC MySQL连接器,2.x版本坐标必须是com.ververica --> <dependency> <groupId>com.ververica</groupId> <artifactId>flink-connector-mysql-cdc</artifactId> <version>${flink.cdc.version}</version> </dependency> </dependencies>

这里解释一下每个依赖的作用:flink-streaming-java提供DataStream API和Source算子,flink-clients负责把任务提交到LocalEnvironment执行,flink-table-api-java-bridge是为了兜底——万一你后续想用CREATE TABLE的方式建CDC源表,不用再改pom。最核心的是flink-connector-mysql-cdc,它内部封装了Debezium引擎、MySQL binlog客户端和Source实现。

2.3 主类代码:MysqlCdcDemo

整个Demo的核心就是一个main方法,代码如下:

package com.example.cdc; import com.ververica.cdc.connectors.mysql.source.MySqlSource; import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class MysqlCdcDemo { public static void main(String[] args) throws Exception { // 1. 创建Flink执行环境,IDE里运行会以Local模式启动 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 2. 强制开启checkpoint,CDC的binlog位点依赖状态做持久化 env.enableCheckpointing(5000); // 3. 构建MySQL CDC Source MySqlSource<String> mySqlSource = MySqlSource.<String>builder() .hostname("localhost") .port(3306) .databaseList("test_db") // 要捕获的数据库,支持正则 .tableList("test_db.user_info") // 要捕获的表,格式:库名.表名 .username("cdc_user") .password("cdc_password") .serverTimeZone("Asia/Shanghai") .deserializer(new JsonDebeziumDeserializationSchema()) .build(); // 4. 接入Source,不加Watermark策略即可 DataStream<String> stream = env.fromSource( mySqlSource, WatermarkStrategy.noWatermarks(), "mysql-cdc-source" ); // 5. 直接打印到控制台,方便观察输出 stream.print().setParallelism(1); // 6. 提交并执行 env.execute("mysql-cdc-demo"); } }

这段代码你只需要改hostname、port、databaseList、tableList、username、password这六个参数,就可以在自己的MySQL上跑起来。

2.4 代码背后的三个关键设计选择

为什么用env.fromSource而不是老的env.addSource?Flink 1.12之后推荐新的Source API,addSource被标记为废弃。fromSource支持Watermark策略、并行度动态调整等新特性,CDC连接器在2.x版本里也全面迁移到了新Source接口。你用老接口虽然也能编译,但会在日志里看到废弃警告,某些版本的CDC还会出现位点提交不一致的问题。

为什么反序列化器选JsonDebeziumDeserializationSchema?CDC Source底层是Debezium引擎,产出的原始数据是结构化的ChangeEvent对象。这个反序列化器把它拍平成一行JSON字符串,包括before、after、op、source等字段,既方便print观察,也方便后续接Kafka或写文件。

为什么必须开checkpoint?这是整套Demo最容易被忽略但最重要的一行。CDC任务启动时会先做一次历史数据快照,然后从binlog的某个位点开始消费增量,这个位点保存在Flink的状态(State)里。不开checkpoint,任务一旦重启,状态全部丢失,Source会重新做快照——重复消费还是小事,生产环境里这会造成严重的数据重复。所以在Demo阶段就养成开checkpoint的习惯,后面上生产才不会踩坑。

3. MySQL端只有三步:binlog、权限、连接参数

很多人在代码里折腾半天,却忘了数据库本身要配合。Flink CDC本质上是伪装成一个MySQL备库去读binlog,MySQL端不给开权限、不开binlog,代码写得再对也没用。

3.1 检查与开启binlog

第一步先确认你的MySQL有没有开启binlog,执行下面的SQL:

SHOW VARIABLES LIKE 'log_bin'; SHOW VARIABLES LIKE 'binlog_format'; SHOW VARIABLES LIKE 'binlog_row_image';

如果log_bin是OFF,需要修改MySQL配置文件(Linux下通常是/etc/my.cnf或/etc/mysql/mysql.conf.d/mysqld.cnf,Windows下是my.ini),在[mysqld]段落下追加:

[mysqld] server-id=1 log_bin=mysql-bin binlog_format=ROW binlog_row_image=FULL expire_logs_days=7

然后重启MySQL服务。改完后重新执行上面的SQL确认三项都满足:log_bin=ON、binlog_format=ROW、binlog_row_image=FULL。

这里补一个原理说明:binlog_format=ROW表示binlog记录的是每一行变更前后的值,而不是SQL语句本身,这是CDC能拿到完整镜像的前提;binlog_row_image=FULL表示记录整行的所有字段,如果设成MINIMAL,UPDATE事件里只包含被修改的字段和主键,before镜像会缺失,Debezium输出的数据就不完整。这两个参数缺一不可。

3.2 创建专用账号与授权

不建议直接拿root账号跑CDC,创建一个专用账号,权限给到最小即可:

-- 创建用户,密码按需修改 CREATE USER 'cdc_user'@'%' IDENTIFIED BY 'cdc_password'; -- 授权 GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'cdc_user'@'%'; -- 刷新权限 FLUSH PRIVILEGES;

每个权限背后都有它的用处,我列个表说明:

权限为什么需要
SELECT读取表的历史数据,用于启动阶段的快照
RELOAD执行FLUSH TABLES WITH READ LOCK获取一致性快照
SHOW DATABASES枚举库列表,用于databaseList过滤
REPLICATION SLAVE读取binlog原始事件
REPLICATION CLIENT查询MySQL主从状态、server-id等元信息

如果你用的是MySQL 8.0,还需要确认一下JDBC驱动版本。MySQL 8.0默认的认证插件是caching_sha2_password,老版本的mysql-connector-java不支持,会报“Authentication plugin”相关错误。Flink CDC 2.4内部自带的驱动版本可以处理这个问题,但如果你的工程里单独引了旧版MySQL驱动,就可能冲突。

3.3 连接参数里的细节

代码里有一组参数容易被忽略,单独拿出来说:

  • databaseList和tableList:前者写库名,后者必须写成库名.表名。多说一句,这两个参数是支持正则的,比如test_db\..*能匹配test_db下所有表,但Java字符串里反斜杠需要转义为\\.新手最容易在这里写错。
  • serverTimeZone:建议显式指定为Asia/Shanghai,不要依赖MySQL服务器时区。Debezium在处理TIMESTAMP类型时会转成UTC存储,如果时区不一致,你看到的数据会差8小时,排查起来非常困惑。
  • server-id:这个参数代码里没写,属于可选配置,但生产环境强烈建议显式设置。Flink CDC会伪装成一个MySQL从库去拉binlog,它需要一个server-id.如果MySQL实例上已经有其他主从复制在跑,随机生成的server-id可能冲突,导致连接瞬间断开。官方默认的规则是从5400开始随机生成,我在下面的坑位章节会详细展开。

4. 跑起来看效果:从快照到增量的完整观测链路

代码和MySQL端都准备好之后,我们来完整走一遍运行流程,同时把输出数据逐字段看清楚——这一节看完,你就知道CDC到底给你吐出了什么东西。

4.1 准备测试数据

先在MySQL里建一个测试库和表:

CREATE DATABASE test_db; USE test_db; CREATE TABLE user_info ( id INT PRIMARY KEY, name VARCHAR(64), age INT, update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); INSERT INTO user_info (id, name, age) VALUES (1, '张三', 20); INSERT INTO user_info (id, name, age) VALUES (2, '李四', 25); INSERT INTO user_info (id, name, age) VALUES (3, '王五', 30);

然后修改代码里的连接信息,跑main方法。

4.2 第一次启动:先打快照

任务启动后,Flink CDC会先扫描test_db.user_info的当前数据,把表里的存量记录全部输出一遍。这就是“快照阶段”。你会在控制台看到类似这样的三条JSON:

{"before":null,"after":{"id":1,"name":"张三","age":20,"update_time":1735000000000},"op":"r","source":{"db":"test_db","table":"user_info","ts_ms":1735000000000}}

关键信息是"op":"r",这里的r代表READ,意思是数据来自启动阶段的快照读取,而不是binlog里的实时变更。同时"before":null表示这是一条插入前的空镜像——快照读没有“前值”。

看到这三条op=r的记录,说明代码和数据库配置已经通了。如果表里有10万条数据,这里会输出10万条JSON,这就是CDC的“全量初始化”能力。

4.3 增量变更实验:insert、update、delete

任务保持运行,回到MySQL依次执行三类DML,观察控制台输出。

执行一条INSERT:

INSERT INTO user_info (id, name, age) VALUES (4, '赵六', 28);

控制台对应输出:

{"before":null,"after":{"id":4,"name":"赵六","age":28,"update_time":1735000100000},"op":"c","source":{"db":"test_db","table":"user_info","lsn":123456,"ts_ms":1735000100000}}

注意op从r变成了c,即CREATE,表示这是一条binlog里的插入事件。

再执行一条UPDATE:

UPDATE user_info SET age = 21 WHERE id = 1;

控制台输出:

{"before":{"id":1,"name":"张三","age":20,"update_time":1735000000000},"after":{"id":1,"name":"张三","age":21,"update_time":1735000200000},"op":"u","source":{...}}

这就能看出before和after的完整语义了:before是变更前的行镜像,after是变更后的行镜像。

最后执行DELETE:

DELETE FROM user_info WHERE id = 2;

对应输出:

{"before":{"id":2,"name":"李四","age":25,"update_time":1735000000000},"after":null,"op":"d","source":{...}}

DELETE事件的after为null,before保留了被删除行的完整数据。

4.4 这里值得单独观察的字段

source对象里有几个字段,source.ts_ms是binlog里事件写入的时间戳(毫秒),source.lsn是binlog文件的位点序号。如果你要做“变更数据落库时间延迟”的监控,ts_ms和时间字段一口径对比就能算出端到端延迟。

另外注意update_time这个字段,代码里的输出是一个微秒级别的数字(例如1735000000000)。这是Debezium对TIMESTAMP类型的默认映射,它返回的是epoch毫秒值,不是字符串。如果你希望直接看到2024-11-28 10:00:00这种格式,要么在SQL层处理,要么自定义反序列化器——这是另一个话题,先记住这个特征。

4.5 验证checkpoint确实在工作

本地跑Demo时看不到Web UI,但可以留意控制台日志。开启checkpoint后,日志里定期会出现类似下面的信息:

Completed checkpoint 2 for job ... (duration: 25ms)

如果你没看到这行日志,说明checkpoint没生效,检查代码里是否执行了env.enableCheckpointing(5000),以及你的环境是否处于Local模式。这一步极容易被忽略,却是整个CDC能“断点续传”的根基。

5. 高频踩坑实录:版本、序列化与数据格式问题

这一节汇总我帮别人排查Flink CDC Demo时遇到的高频问题。每一个都是我亲眼见过的,按“现象-原因-修复”的顺序写。

5.1 编译失败:找不到类MysqlSource / MySqlSource

这个报错信息通常是:

Cannot resolve symbol 'MysqlSource'

或者:

Cannot resolve symbol 'MySqlSource'

原因很简单:Flink CDC 1.x时代的类名是MysqlSource(注意中间是小写y),内部实现基于废弃的SourceFunction;从2.0开始改成了MySqlSource(中间大写S),底层换成了新的Source API。你从老博客复制代码,类名就会对不上。

解决方案一句话:认准MySqlSource(M大写、y小写、S大写)。如果你看到的教程里写的是MysqlSource,大概率是2021年之前的文章,请直接关掉。

5.2 坐标混用:com.ververica 还是 org.apache.flink

如果你的pom.xml里写的是:

<groupId>org.apache.flink</groupId> <artifactId>flink-cdc-connector-mysql</artifactId>

但代码里用的是2.x的com.ververica.cdc.connectors.mysql.source.MySqlSource,编译必然报错。反过来,如果坐标是com.ververica但你的Flink版本是1.15以下,也可能出现工厂加载失败。

我的建议:这套Demo固定用com.ververica:flink-connector-mysql-cdc:2.4.2,不要边抄边改,先把库跑通再去升级。

5.3 运行报错:Could not find any factory for identifier 'mysql'

完整报错类似:

Could not find any factory for identifier 'mysql' that implements 'org.apache.flink.table.factories.DynamicTableFactory'.

这个报错通常出现在你用CREATE TABLE ... WITH ('connector' = 'mysql-cdc')的SQL方式接入CDC时。原因有三种:

一是flink-connector-mysql-cdc的jar没有进classpath,mvn dependency:tree看一眼就知道。

二是Flink版本和连接器版本不匹配。Flink 1.17配CDC 2.4.2没问题,但如果Flink是1.13,配CDC 2.4就大概率找不到工厂,因为连接器2.4的最低要求是Flink 1.15。

三是你的工程里只引了flink-table-api-java,没有引入flink-table-planner-loader。本地IDE跑SQL方式时,请确认pom里至少有一个planner依赖。本Demo直接用DataStream API的fromSource,绕开了工厂机制,所以反而不会遇到这个报错——这也是我推荐先用DataStream API的原因之一。

5.4 只有历史数据,后续增删改没有反应

这是最让人头疼的现象:启动后打出了已有数据,但执行INSERT/UPDATE/DELETE,控制台无动于衷。按顺序排查:

第一步,检查binlog是否真的开启,尤其是MySQL 8.0,默认binlog是开启的,但5.7的某些发行版默认是关闭的,执行SHOW VARIABLES LIKE 'log_bin'确认。

第二步,检查binlog_format是否为ROW。有些云数据库默认是STATEMENT或MIXED,Debezium对格式有严格要求。修改后记得重启MySQL。

第三步,检查账号权限。如果缺少REPLICATION SLAVE权限,通常会直接报Access denied,但如果你用的账号其实是root,那这一步通常没有问题。

第四步,确认代码里的databaseList和tableList确实匹配你的表和库名。tableList写test_db.user_info就只监听这一张表,如果执行INSERT时写错了库名,任务自然没反应。

第五步,也是最隐蔽的:MySQL的binlog开启时间晚于表的创建时间。如果这个表在binlog开启之前就存在,任务启动时快照阶段读取的正常,但增量阶段要从binlog的某个位点开始找,如果这个位点之前的binlog已经被清理掉,任务会卡住或者静默不消费。解决办法就是给binlog保留更长时间,或者在开启binlog后重新创建测试表。

5.5 server-id冲突:连接一会儿就断开

如果运行日志里出现:

Something unusual happened: MySQL connection was killed

大概率是server-id冲突。你的MySQL实例可能已经有一个主从复制在使用某个server-id,Flink CDC随机生成的server-id恰好撞上了。解决办法很简单,在builder里显式指定:

.serverId("5400-5404")

这段配置表示给这个Source分配5400到5404这5个server-id,映射关系是并行度决定的——单并行度用5400,双并行度用5400和5401,以此类推。分配区间避开你现有主从的server-id即可。

5.6 时间字段变成微秒数字,怎么读

很多人在输出里看到update_time: 1735000000000这种字段时都会懵一下。Debezium对MySQL的TIMESTAMP类型默认转成epoch毫秒值,这个设计是为了避免时区歧义。如果你需要在控制台看到可读格式,最简单的办法是换一个自定义反序列化器,或者干脆用SimpleStringSchema接收原始字节再自己解析。

这里我插一句个人建议:如果你只是想快速验证数据通道通没通,完全不用纠结这个格式,数字时间戳一样能证明数据在流动。但如果你要写业务逻辑,建议考虑用Table API定义CDC表而不是直接消费JSON字符串,Table API下Flink会按照MySQL元数据把类型映射成正确的SQL类型,时间字段直接就是时间类型。

5.7 从print到JDBC Sink的扩展路径

最后说一个高频需求:跑通print之后,很多人的下一步是把数据写回另一个MySQL库或者Kafka。如果写MySQL,需要在pom里额外引入flink-connector-jdbc,然后自己实现一个Sink,处理op=c/u/d三种事件对应的INSERT/UPDATE/DELETE语句。

这里有个循环陷阱要提醒:如果目标表和源表在同一个MySQL实例,CDC任务会把自己写入的数据也当成binlog事件再读回来,形成无限循环。生产环境的做法一定是分实例,或者写目标表时用op字段过滤掉回环数据。这个点我在实际项目中见过不止一次,写下来给你提个醒。

最后再分享一个实操上的小事:这套Demo里stream.print().setParallelism(1)这行,并行度设成1是为了让你看日志时不乱序。如果改成大于1,多并行度下打印顺序会交叉,排查问题容易看花眼。等我需要压测或者真实接入业务时,记得把并行度提上去,同时把serverId的范围扩展成和并行度匹配的数量——这两个参数是绑定的。

这套组合我把Flink 1.17、Flink-CDC 2.4、Java 11固定成了自己的默认模板,配合MySQL 5.7和8.0都跑得很稳定。下一个阶段你可以尝试把print替换成JDBC Sink写回业务库,或者接一个Kafka Topic做数据分发,但无论往哪个方向走,binlog的配置和checkpoint 的逻辑都是不变的地基。先把地基打牢,后面的事情会顺很多。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询