提个问题你在网上搜“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.flinkFlink CDC的项目坐标在2.x和3.x之间发生过一次“搬家”。2.4及之前的版本groupId是com.ververicaartifactId是flink-connector-mysql-cdc到了3.0之后整个项目进入Apache Flink组织groupId变成了org.apache.flinkartifactId也改成了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.17MySQL5.7或8.0必须开启binlog且格式为ROWJDK8或11建议11别用17去跑Flink 1.17Maven3.6工程构建工具没什么特殊要求这套组合最大的优势是“文档密度高”。Flink 1.17是社区使用最广的版本之一Flink CDC 2.4的API设计又保留了最直观的DataSource写法你在网上搜到的大部分报错和解决方案都能直接对应上对新手非常友好。1.3 为什么暂时不用Flink-CDC 3.0Flink 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.xmlproperties flink.version1.17.2/flink.version flink.cdc.version2.4.2/flink.cdc.version maven.compiler.source11/maven.compiler.source maven.compiler.target11/maven.compiler.target project.build.sourceEncodingUTF-8/project.build.sourceEncoding /properties dependencies !-- Flink DataStream API -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version${flink.version}/version /dependency !-- 本地运行所需提供LocalEnvironment入口 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version${flink.version}/version /dependency !-- 如果后续想用SQL方式建CDC表需要table api的bridge -- dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge/artifactId version${flink.version}/version /dependency !-- Flink CDC MySQL连接器2.x版本坐标必须是com.ververica -- dependency groupIdcom.ververica/groupId artifactIdflink-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. 强制开启checkpointCDC的binlog位点依赖状态做持久化 env.enableCheckpointing(5000); // 3. 构建MySQL CDC Source MySqlSourceString mySqlSource MySqlSource.Stringbuilder() .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策略即可 DataStreamString 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.addSourceFlink 1.12之后推荐新的Source APIaddSource被标记为废弃。fromSource支持Watermark策略、并行度动态调整等新特性CDC连接器在2.x版本里也全面迁移到了新Source接口。你用老接口虽然也能编译但会在日志里看到废弃警告某些版本的CDC还会出现位点提交不一致的问题。为什么反序列化器选JsonDebeziumDeserializationSchemaCDC 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备库去读binlogMySQL端不给开权限、不开binlog代码写得再对也没用。3.1 检查与开启binlog第一步先确认你的MySQL有没有开启binlog执行下面的SQLSHOW 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.cnfWindows下是my.ini在[mysqld]段落下追加[mysqld] server-id1 log_binmysql-bin binlog_formatROW binlog_row_imageFULL expire_logs_days7然后重启MySQL服务。改完后重新执行上面的SQL确认三项都满足log_binON、binlog_formatROW、binlog_row_imageFULL。这里补一个原理说明binlog_formatROW表示binlog记录的是每一行变更前后的值而不是SQL语句本身这是CDC能拿到完整镜像的前提binlog_row_imageFULL表示记录整行的所有字段如果设成MINIMALUPDATE事件里只包含被修改的字段和主键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表示这是一条插入前的空镜像——快照读没有“前值”。看到这三条opr的记录说明代码和数据库配置已经通了。如果表里有10万条数据这里会输出10万条JSON这就是CDC的“全量初始化”能力。4.3 增量变更实验insert、update、delete任务保持运行回到MySQL依次执行三类DML观察控制台输出。执行一条INSERTINSERT 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里的插入事件。再执行一条UPDATEUPDATE 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是变更后的行镜像。最后执行DELETEDELETE FROM user_info WHERE id 2;对应输出{before:{id:2,name:李四,age:25,update_time:1735000000000},after:null,op:d,source:{...}}DELETE事件的after为nullbefore保留了被删除行的完整数据。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。你从老博客复制代码类名就会对不上。解决方案一句话认准MySqlSourceM大写、y小写、S大写。如果你看到的教程里写的是MysqlSource大概率是2021年之前的文章请直接关掉。5.2 坐标混用com.ververica 还是 org.apache.flink如果你的pom.xml里写的是groupIdorg.apache.flink/groupId artifactIdflink-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没有进classpathmvn 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或MIXEDDebezium对格式有严格要求。修改后记得重启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-idFlink 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处理opc/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 的逻辑都是不变的地基。先把地基打牢后面的事情会顺很多。