资讯动态

Flink CDC实战:MySQL实时同步到Kafka完整指南

发布时间:2026/9/16 23:23:12 来源:尧图企业网站定制
做实时数仓和实时同步的兄弟应该都绕不过一件事怎么把MySQL里的业务变更数据实时搬到其他系统里去。我这边经历过好几轮方案选型从最开始的Canal客户端消费到后来定时刷增量最后稳定跑在Flink CDC Kafka这条链路上。今天就把DataStream方式的完整做法、关键参数和踩过的坑一次说清楚照着敲基本能跑通。这篇文章适合两类人看一类是刚接触实时同步、想快速搭一套MySQL到Kafka链路的同学另一类是从Flink SQL方式切过来、发现DataStream方式更灵活可控的开发者。我会把从环境准备、依赖版本、完整代码到问题排查的整个链路都展开讲尽量少说空话多给能直接落地的细节。1. 链路设计与方案拆解1.1 这套链路到底解决了什么问题先聊一下背景。业务库的MySQL每天都在产生增删改下游的ES、Redis缓存、数仓ODS层都需要感知这些变化。过去最常见的做法是业务代码里双写或者定时任务扫增量字段但这两种方式都有明显毛病双写侵入业务代码、耦合重定时扫表延迟高、对线上库有压力而且删操作不好捕获。Flink CDC做的事很简单直接读MySQL的binlog把每一步INSERT、UPDATE、DELETE都解析成结构化事件再原样或者加工后发出去。搭配Kafka做消息管道下游所有系统只要订阅Kafka就能拿到同一份变更数据做到了业务无侵入、准实时、一处采集多处消费。这个模式现在几乎是实时数仓的标配特别是ODS层同步基本都是这条路子。我之所以强调DataStream方式而不是Flink SQL方式是因为DataStream方式更接近底层你能控制每个环节的序列化器、并行度、投递语义也能在中间任意插入Map、Filter、异步IO做加工。Flink SQL方式写起来虽然短但遇到自定义消息格式、动态topic、复杂维表关联的时候反而束手束脚。如果目标是把CDC数据原样推到Kafka、后续再让其他程序消费DataStream是最直接也最好排查问题的写法。1.2 DataStream、SQL两种方式怎么选很多教程一上来就讲Flink SQL但实际项目里我更推荐先理解DataStream。Flink CDC底层的Source和Debezium解析逻辑是一样的两种方式都是基于同一套增量快照框架区别在于你拿到的“数据形态”和“可编程性”。用Flink SQL你写一条CREATE TABLE INSERT INTO SELECT就完事但它内部会帮你做类型映射、把Debezium的changelog转成Flink内部的RowData一旦你想拿到原始的before/after完整JSON反而要额外设置。DataStream方式直接用JsonDebeziumDeserializationSchema拿到的就是Debezium格式的JSON字符串原汁原味后面想怎么处理都行。这里给个经验值如果只是“MySQL binlog整表同步到Kafka topic”两个方式都能干但DataStream更适合需要控制消息key、需要按表拆topic、需要自定义序列化的场景。而如果后面要直接接Flink SQL做实时JOIN、聚合那SQL方式的一体化体验更好。我个人的习惯是采集层统一用DataStream进Kafka计算层再用SQL从Kafka读各司其职维护成本最低。1.3 整体数据流长什么样这条链路从数据流上看其实很短MySQL binlog → Flink CDC Source → Kafka Sink。Flink CDC Source内部会先做一次一致性快照把当前全量数据读出来同时记录当时的binlog位点快照完成后无缝切换到实时增量继续解析binlog事件。整个过程对业务侧透明不需要MySQL重启也不需要业务改代码。数据到了Kafka之后消息体是Debezium的JSON格式包含before、after、op、source等信息。之后ES同步、Redis缓存更新、数仓入湖都是下游消费者各自的事与采集端解耦。这套设计的核心收益就是采集链路一旦跑稳定后续加新的下游消费者几乎零成本只要订阅同一个topic即可。2. 前置准备MySQL、Kafka、Flink环境三件套2.1 MySQL端binlog必须开成ROW格式Flink CDC读的是binlog如果MySQL没开binlog或者binlog格式不是ROW那Source连上之后什么都拿不到或者拿到的数据不完整。这里先强调一个关键概念binlog有三种格式STATEMENT、ROW、MIXED。CDC必须用ROW格式因为只有ROW格式才记录了每行数据变更前后的完整值STATEMENT只记录SQL语句解析不出before/after。在MySQL的配置文件里加上这几行然后重启MySQL[mysqld] server-id 1 log_bin mysql-bin binlog_format ROW binlog_row_image FULL expire_logs_days 7解释一下这几个参数server-id在MySQL集群里必须唯一Flink CDC也会作为一台“假从库”连上来拉binlogbinlog_row_image设置为FULL是为了让日志里同时保留变更前和变更后的完整行镜像如果设置成MINIMAL有些字段就不会记录你拿到的before可能是空的下游做对比会出问题。expire_logs_days控制binlog保留天数建议至少保留3天以上防止Flink任务停机恢复时binlog已经被清理。提示MySQL 8.0里的binlog_expire_logs_seconds可以替代expire_logs_days两个任选一个就行不要同时配置导致启动报错。2.2 给CDC账号开最小必要权限很多朋友第一步就卡在权限上明明配了账号密码Flink任务一启动就报Access denied。用root虽然能跑通但生产环境不推荐按最小权限原则单独建一个CDC账号CREATE USER cdc% IDENTIFIED BY cdc123456; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc%; FLUSH PRIVILEGES;这几个权限分别干什么我简单说下SELECT是快照阶段查全量数据用的REPLICATION SLAVE和REPLICATION CLIENT是拉取binlog和查看binlog位点必需的RELOAD用于快照时获取一致性锁SHOW DATABASES是让Source能枚举库表。如果用的是MySQL 8.0的默认认证插件caching_sha2_password某些版本连接会报错可以在创建用户时指定mysql_native_password或者直接用支持该认证的连接器版本。2.3 依赖版本直接决定你能不能跑起来版本问题在Flink CDC里特别坑因为Flink、CDC连接器、Kafka连接器三者之间存在严格的兼容矩阵。我目前线上稳定跑的组合是Flink 1.17.2 Flink CDC 2.4.2 Kafka连接器1.17.2 Kafka 3.4。这个组合经过大量验证网上资料也最多出问题好查。Flink CDC 3.x现在也出了架构变化比较大支持了Pipeline方式但如果刚上手我建议先用2.4.x把链路跑通再考虑升级。Maven依赖这样配dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version1.17.2/version /dependency dependency groupIdcom.ververica/groupId artifactIdflink-connector-mysql-cdc/artifactId version2.4.2/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version1.17.2/version /dependency注意flink-connector-kafka的版本号跟Flink主版本保持一致不是跟Kafka集群版本保持一致。你只要用这个版本对应的连接器它内置的Kafka客户端会自己跟集群兼容不需要额外引入kafka-clients依赖。另外如果你打包的时候发现flink-streaming-java带scope为provided本地IDE运行时记得去掉或者改成compile否则会报找不到类。3. 核心实现MySQL CDC Source与Kafka Sink完整代码3.1 构建MySQL CDC Source的几个核心参数MySQL CDC Source一般用MySqlSource.builder()来构建核心参数就那么几个但每个都值得掰开说。hostname和port指向MySQL地址databaseList和tableList决定监听范围可以填正则比如databaseList(test_db)配合tableList(test_db.users)也可以直接监听整个库。username和password就是刚才建的CDC账号。这里有个细节参数没有填的字段走的是默认值而默认值不一定适合你的场景。比如serverTimeZone默认是UTC如果业务库是北京时间而且时间字段是DATETIME类型不配会差8个小时。我建议显式配置serverTimeZone(Asia/Shanghai)。startupOptions是另一个关键参数它决定任务启动时从哪个位置开始读选项行为适用场景initial()先做全量快照再无缝切增量首次接入、目标端没有历史数据latest()跳过全量只读启动之后的新增变更目标端已有存量数据只要增量earliest()从最早可用的binlog开始需要回放历史变更specificOffset()从指定的binlog文件名位点开始断点续传、故障恢复这几个选项选错的风险我遇到过一次某次任务重启直接配了latest()结果重启间隙漏了几分钟数据下游对不上账排查很久才发现是启动模式的问题。所以一般建议首次同步用initial()之后运维重启不要随便改配置。3.2 反序列化器决定Kafka里消息长什么样deserializer这一步直接决定你发到Kafka里的消息格式。最省事的是JsonDebeziumDeserializationSchema它把Debezium的变更事件转成JSON字符串结构大概是这样{ before: null, after: { id: 1, name: 张三, age: 25 }, source: { version: 1.9.7.Final, connector: mysql, name: mysql_binlog_source, db: test_db, table: users, server_id: 223344, ts_ms: 1710000000000 }, op: c, ts_ms: 1710000001234 }op字段表示操作类型c是新增u是更新d是删除r是快照阶段读取的全量数据。before和after分别是变更前后的行数据删除时after为null新增时before为null。source里带库名、表名、server_id和binlog时间戳下游可以根据source.db和source.table判断数据来自哪张表。如果你想自定义消息格式比如只要after部分、或者把更新和删除统一包装成自己的JSON结构可以继承DebeziumDeserializationSchema自己写。实际项目中我经常这么干因为直接给下游原始Debezium格式他们还得自己解析不如在采集端就统一成DTO结构但缺点是采集端耦合了下游格式看团队取舍。3.3 Kafka Sink的投递语义与参数配置Flink 1.15之后官方推荐用KafkaSink替代了老的FlinkKafkaProducer。构建KafkaSink最简单的方式是设置bootstrap地址和topic再指定value的序列化器KafkaSinkString kafkaSink KafkaSink.Stringbuilder() .setBootstrapServers(127.0.0.1:9092) .setRecordSerializer( KafkaRecordSerializationSchema.builder() .setTopic(mysql-cdc-users) .setValueSerializationSchema(new SimpleStringSchema()) .build() ) .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE) .build();这里最重要的概念是DeliveryGuarantee也就是投递语义。AT_LEAST_ONCE表示消息至少投递一次极端情况下可能重复EXACTLY_ONCE表示精确一次需要配合Flink的Checkpoint开启两阶段提交。CDC场景我一般建议先上AT_LEAST_ONCE因为数据本身带主键和op类型下游做幂等很容易。如果你真的需要EXACTLY_ONCE那必须保证Kafka broker的transaction.max.timeout.ms不小于Flink的transaction.timeout否则启动会报错这个后面排查章节详细说。另外补充一个容易被忽略的点KafkaSink默认会把消息均匀分布到topic的各个分区如果对同一行数据的顺序有要求必须设置消息key。cdc场景下同一主键的变更最好保证进同一个分区否则下游按主键做状态更新会乱序。可以通过setKeySerializationSchema配合提取主键字段来设置key我习惯把主键值或者表名主键值拼成key。3.4 完整可运行的示例工程代码下面这段是完整的可运行代码环境是Flink 1.17 Flink CDC 2.4.2我已经把Source、Sink、Checkpoint都配置好了直接复制改参数就能跑import com.ververica.cdc.connectors.mysql.source.MySqlSource; import com.ververica.cdc.connectors.mysql.table.StartupOptions; import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.connector.base.DeliveryGuarantee; import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema; import org.apache.flink.connector.kafka.sink.KafkaSink; import org.apache.flink.runtime.state.hashmap.HashMapStateBackend; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.environment.CheckpointConfig; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class MySqlCdcToKafkaJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000L); env.setStateBackend(new HashMapStateBackend()); env.getCheckpointConfig().setCheckpointStorage(file:///data/flink/checkpoints); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(1000L); env.getCheckpointConfig().setCheckpointTimeout(60000L); env.getCheckpointConfig() .setExternalizedCheckpointCleanup( CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION); MySqlSourceString mySqlSource MySqlSource.Stringbuilder() .hostname(127.0.0.1) .port(3306) .databaseList(test_db) .tableList(test_db.users) .username(cdc) .password(cdc123456) .serverTimeZone(Asia/Shanghai) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build(); KafkaSinkString kafkaSink KafkaSink.Stringbuilder() .setBootstrapServers(127.0.0.1:9092) .setRecordSerializer( KafkaRecordSerializationSchema.builder() .setTopic(mysql-cdc-users) .setValueSerializationSchema(new SimpleStringSchema()) .build()) .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE) .build(); DataStreamSourceString stream env.fromSource( mySqlSource, WatermarkStrategy.noWatermarks(), MySQL CDC Source); stream.sinkTo(kafkaSink); env.execute(mysql-cdc-to-kafka); } }这段代码里有个容易被忽视的设计WatermarkStrategy.noWatermarks()。因为CDC Source处理的不是无界流式事件时间它本身不依赖水印做窗口计算直接用noWatermarks避免无意义的watermark生成开销。如果你想在Source和Sink之间做基于事件时间的窗口聚合那才需要配置水印但现在这个场景不需要。还有一个经验本地IDE直接跑这个main方法是可以的但打包提交到Flink集群时注意不要把flink-streaming-java打进去否则会和集群自带的冲突。用maven-shade插件打包时把scopeprovided的依赖排除掉或者直接在pom里设置provided就能避免ClassNotFound或者冲突异常。3.5 Checkpoint配置为什么是重中之重做CDC同步Checkpoint不是可选项是必选项。Flink CDC的增量机制依赖Checkpoint保存binlog位点任务重启后从最近一次Checkpoint恢复才能做到不丢数据。上面代码里env.enableCheckpointing(5000L)表示每5秒做一次Checkpoint这个间隔可以根据业务容忍的恢复延迟调整间隔越短故障恢复时丢失的数据越少但对磁盘和状态的IO压力也越大。CheckpointStorage建议配上文件系统的路径比如HDFS或者本地文件系统不要用默认的JobManager内存存储否则Checkpoint一多内存就爆。生产上我一般用HDFS路径本地测试就写file:///data/flink/checkpoints。另外RETAIN_ON_CANCELLATION这个配置也很重要它表示手动取消任务时保留Checkpoint这样你可以从保存的Checkpoint手动恢复跳过重新全量快照。我踩过的一个教训是刚搭环境时没开Checkpoint任务跑了几天一次网络抖动导致JobManager重启重启后Flink CDC又从initial()开始做全量快照把几百万条数据重新发了一遍Kafka里全是重复消息。从那以后我要求的底线是CDC采集任务必须开Checkpoint而且启动方式必须指定从Checkpoint恢复。4. 实操验证启动任务在MySQL增删改去Kafka看数据4.1 打包提交任务前的准备代码写完之后需要先准备Kafka的topic。这里有个常见问题Kafka生产消费命令启动一次会一直运行吗如果你用的是kafka-console-producer.sh和kafka-console-consumer.sh这类命令行工具只要不按CtrlC进程就会一直挂着等待输入或持续消费所以在测试环境跑通了记得手动退出不然终端一直占着。先创建topic分区数建议跟下游消费并行度匹配我这里建3个分区kafka-topics.sh --bootstrap-server 127.0.0.1:9092 \ --create --topic mysql-cdc-users \ --partitions 3 --replication-factor 1然后启动一个消费者在后台挂着方便等下直接看数据kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \ --topic mysql-cdc-users --from-beginning打包提交命令我就不重点讲了Flink集群用flink run -c 主类名 jar包路径即可。本地IDE直接跑main方法也完全可行因为Flink会以local模式启动一个嵌入式集群。4.2 在MySQL执行各种DML观察Kafka里的消息假设test_db.users表里原本有三条数据CREATE TABLE users ( id INT PRIMARY KEY, name VARCHAR(50), age INT ); INSERT INTO users VALUES (1, 张三, 25), (2, 李四, 30), (3, 王五, 28);任务用initial()方式启动后你会先在Kafka里看到三条op为r的快照消息每条消息对应一张表里的一行after里就是完整的行数据。快照读完之后再执行任何DMLKafka里就会实时出现对应的op消息UPDATE users SET age 26 WHERE id 1; DELETE FROM users WHERE id 3; INSERT INTO users VALUES (4, 赵六, 22);如果你眼睛够快能看到update和delete几乎在一条事务提交后立刻出现在Kafka里延迟通常在毫秒到秒级主要取决于binlog的推送频率和Kafka的写入耗时。update消息里before是变更前的旧值after是新值这对下游做“变更前对比”非常有用很多数仓的拉链表就是靠这个实现的。4.3 从Kafka消费到的消息怎么排错如果Kafka消费者里半天没动静先不要怀疑代码按下面顺序排查先看topic是否存在再看Kafka集群是否正常再看Flink任务日志里有没有报错。很多时候Flink任务其实已经在报错了只是你没有看日志导致误以为还在运行。用kafka-console-consumer看到的消息如果格式不对比如乱码或者缺字段优先检查deserializer。JsonDebeziumDeserializationSchema输出的是标准JSON字符串如果Kafka消费者里看到了类似{before:null,after:{...}}的结构说明链路是通的。如果你自定义了序列化器排错时可以先临时换回JsonDebeziumDeserializationSchema确认问题出在Source还是Sink这个二分法能帮你省很多时间。5. 常见问题与排查实录5.1 权限、binlog配置不对导致连接失败最典型的报错是Communications link failure或者Access denied for user。前者一般是网络不通、MySQL没开外网访问、或者server-id冲突后者多半是权限没给全。我整理了一份问题速查表方便直接对照报错现象大概率原因解决办法Access denied for userCDC账号权限不足按上文GRANT语句重新授权The server is not configured to use binlogbinlog没开启检查my.cnf并重启MySQLbinlog format must be ROWbinlog格式不对设置binlog_formatROWFound conflicting server idMySQL实例自己配了server_id给MySQL和Flink CDC分配不同的server_idTable xxx doesnt existtableList填错检查库名表名大小写、通配符写法Unrecognized MySQL connector version连接器版本与MySQL版本不匹配升级Flink CDC连接器版本其中一个我特别想强调如果MySQL是云厂商的托管实例比如RDS部分云厂商默认隐藏了binlog的某些权限你即使开了binlog也可能拉不到。这时候要去控制台确认binlog是否开放、binlog保留时间是否够长。有些云环境还会限制REPLICATION权限遇到这种情况只能联系客服或者改用云平台自带的DTS同步服务。5.2 Kafka消费不到数据Source却显示正常这种问题最磨人Flink任务状态是RUNNINGCheckpoint也正常MySQL有变更但Kafka里就是没消息。先排除一个低级错误——topic写错了或者消费者订阅的和Sink写入的不是同一个topic这个我身边真有人犯过。再就是KafkaSink的setTopic是静态的消息会全部写入这个topic如果topic不存在且broker配置的auto.create.topics.enable为falseSink会一直重试表现在任务日志里是一堆TimeoutException。还有一种隐蔽情况Kafka集群是SASL或SSL认证的你本地测试用明文连但Flink任务提交到集群后走的是另一个网络环境认证信息没配。KafkaSink的构建方法里可以setProperty(security.protocol, SASL_PLAINTEXT)等方式传入认证参数生产环境务必核对。如果以上都没问题再去看topic的offset和消费者组的lag。用kafka-consumer-groups.sh查看消费组的当前offset和log-end-offset能快速判断消息是根本没生产出来还是生产出来了没被你的消费者消费掉。这个命令在排查Kafka链路问题时几乎是必用的建议熟练。5.3 投递语义、重复数据与顺序问题用AT_LEAST_ONCE模式时消息重复是正常现象不是bug。Flink任务重启、网络超时重试、Kafka broker端重试都可能导致下游收到重复消息。应对方式有两种一是下游消费端做幂等按主键去重二是升级到EXACTLY_ONCE但代价是性能下降和配置复杂度上升。EXACTLY_ONCE有个经典坑Flink侧配置事务超时时间默认是1小时而Kafka broker的transaction.max.timeout.ms默认是15分钟两边对不上任务启动后写第一条消息就会报InvalidTxnStateException。解决办法是在KafkaSink里设置transaction.timeout.ms小于等于broker的上限比如KafkaSink.Stringbuilder() .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setProperty(transaction.timeout.ms, 600000)另一个和顺序有关的坑如果不同分区并行处理同一张表的数据而消息没有设置key那么同一行数据的多次更新会分布到不同分区下游按时间顺序回放时会乱序。解决方式我刚才提过用主键做消息key就能保证同一主键的消息永远走同一个分区。这里我自己的经验是如果要保证同一行的严格顺序同时还要做表级顺序那就用“库名.表名.主键”拼key几乎不会踩坑。5.4 性能优化与常见异常处理任务跑到大数据量阶段最常见的两个问题一是全量快照阶段比较慢二是增量阶段背压高。全量快照慢原因往往是Source的并行度是1Flink CDC 2.x的快照读取阶段默认单并行只能靠chunkSize参数来调节每个分片的大小比如MySqlSource.Stringbuilder() .hostname(127.0.0.1) // 其他参数省略 .chunkSize(4096) .build();chunkSize调小TaskManager内存压力小但分片数量多、切换频繁调大单次读的数据多但内存占用高。如果表数据量是千万级建议chunkSize设置8192左右同时给TaskManager足够的堆内存。增量阶段的背压多半是Kafka写入慢或者下游消费慢导致Kafka堆积。这时候去Kafka监控面板看topic的bytes-in和bytes-out如果bytes-in远大于bytes-out说明下游消费跟不上。还有一个容易被忽略的点KafkaSink默认并行度等于TaskManager的slot数如果下游topic的分区数小于Sink并行度会出现部分子任务一直空转或者写冲突。建议Sink并行度和topic分区数一致或者让Sink并行度不大于分区数。Kafka OOM问题也常有人问。Kafka本身是Java进程OOM一般发生在broker端或者消费者端。如果是broker端heap不足调整KAFKA_HEAP_OPTS如果只是消费者拉取超大消息那是消息体本身太大要回头查CDC序列化是不是把整行大字段都塞进去了。对于binlog里的大字段比如TEXT、BLOB下游如果不需要可以在反序列化阶段就裁剪掉别让消息体无限膨胀。6. 踩坑后的几点体会这套链路我在线上跑了一年多最大的体会是越简单的架构越稳能不加中间加工就别加。刚开始我也试过在Flink里做各种字段映射和清洗后来发现下游需求变来变去清洗逻辑放在采集端反而成了改动的瓶颈最后还是老老实实把原始Debezium JSON推到Kafka让各个下游自己解析。还有一点想提醒大家Flink CDC的版本升级一定要谨慎尤其是大版本。Flink CDC 3.x改变了Source的架构和参数升级不是改个版本号那么简单需要重新做回归测试。我的建议是先把当前版本用熟真正理解binlog位点、Checkpoint、快照机制这些底层概念再去看新版本值不值得升级。最后分享一个实用技巧每次修改CDC任务配置前先给Kafka topic设置好消息保留时间比如retention.ms配置成7天。这样即使下游服务挂了几天恢复后还能从Kafka把数据追回来不会因为binlog过期而永久丢失数据。CDC链路的核心是可靠性把这些兜底措施做扎实比追求单点性能更值得花时间。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价