
之前有个需求线上订单表的数据变更要尽量实时同步给数仓和分析团队一开始用定时任务扫update_time延迟能扛可等到业务方要求把“删除”操作也记录下来、还要重新回放历史数据的时候这套方案彻底撑不住了。最后换成了 FlinkCDC 这条链路MySQL 的 binlog 被 FlinkCDC 消费再通过 DataStream API 把变更数据写成 JSON 丢到 Kafka。说白了FlinkCDC 就是一个贴在 MySQL 上监听 binlog 的订阅端Kafka 只是它的下游出口。这篇文章不抄官方文档是我从零跑通这条链路并且线上稳定跑了几个月的完整实操记录适合刚接触实时数据同步、准备用 DataStream 方式接 Kafka 的读者参考。1. FlinkCDC 到底解决了什么问题从一轮轮定时任务说起1.1 从 binlog 说起MySQL 的 binlog二进制日志是 MySQL 在 server 层维护的一份“流水账”里面记录了所有产生数据变更的操作包括 INSERT、UPDATE、DELETE甚至 DDL。MySQL 提供了三种 binlog 格式STATEMENT、ROW、MIXED。FlinkCDC 要求必须用ROW 格式因为只有这种格式会写清楚“哪一行被改成了什么样子”UPDATE 事件里甚至会同时记录这一行更新前的值before和更新后的值after。STATEMENT 格式只记录 SQL 语句本身虽然日志量小但重放时没法保证一致性更没法让下游拿到干净的变更数据。把 binlog 理解成超市的收银小票会更直观每次有人结账收银系统都会打出一张小票记录买了什么、付了多少钱、什么时候发生的。FlinkCDC 做的事情就是专门坐在收银台旁边把小票一张张复印出来再转发给下游。只要这张“小票”持续产生下游就能跟着实时变化。1.2 FlinkCDC 和 Canal、Debezium 有什么区别做 MySQL 实时同步市面上常见的还有 Canal 和 Debezium我也都试过简单做个对比。方案部署方式与 Flink 的集成度适合的场景Canal需要独立部署 Canal Server客户端再订阅不算紧密需要自行封装接入已有 Canal 基础设施、团队熟悉阿里系组件Debezium可以内嵌到 Java 应用也可以配合 Kafka Connect本身是标准 CDC 框架但接 Flink 还要自己处理想用 Kafka Connect 做独立同步管道FlinkCDC直接作为 Flink 的 Source嵌入作业内部天然融合checkpoint、重启策略、watermark 等都能复用希望用一套流处理框架搞定同步和加工我最后选 FlinkCDC 的原因很简单少维护一个中间件。Canal 要单独起一个 server 进程还要考虑它的高可用FlinkCDC 直接作为 Flink 作业里的 Source 运行作业挂了checkpoint 恢复时它会自动从上次记录的位置继续读 binlog不用额外处理。1.3 什么场景别用它FlinkCDC 不是银弹。我遇到一个内部统计系统表数据量不大分钟级延迟完全能接受就没必要引入整套 Flink 和 Kafka维持定时任务反而更省事。FlinkCDC 适合的是数据变更频繁、下游需要实时响应、并且对“一条都不能丢、不能重”有要求的场景。如果只是每天凌晨同步几张表别上这套东西纯属给自己找运维负担。2. 版本兼容是第一道暗坑我的完整选型清单2.1 版本之间到底有多容易出问题FlinkCDC 这种组件版本坑比功能坑更致命。Flink 版本、FlinkCDC 版本、Kafka connector 版本之间没有完全统一的兼容矩阵官方文档写得又比较笼统很多人上来就一把梭最新版结果跑起来编译报错、运行时类冲突、CDC 连接器找不到方法。我最终采用的版本组合如下组件版本说明Flink1.18.1主流稳定版本社区资料多Flink CDC2.4.22.x 系列里比较成熟的版本API 稳定flink-connector-kafka1.18.0必须和 Flink 主版本对应MySQL8.0.33开启 binlogROW 格式Kafka3.4.0用 KRaft 模式跑在 Docker 里方便验证这里有个容易忽略的点flink-connector-kafka的版本号是跟着 Flink 走的不是跟着 Kafka 走的。你 Flink 用 1.18就找对应的flink-connector-kafka-1.18.0。而 FlinkCDC 的版本独立2.4.2 已经验证过能和 Flink 1.18 一起工作。如果你用 Flink 1.14 配 FlinkCDC 2.4大概率会遇到类冲突因为内部依赖的 Flink API 变了。2.2 为什么我坚持用 Docker 搭 Kafka 和 MySQL本地开发阶段最省心的是用 Docker 把依赖的环境一次性拉起来。不用污染本机版本切换也方便。Kafka 3.4 已经支持 KRaft 模式不需要单独再起 ZooKeeperdocker-compose 文件可以精简不少。version: 3.8 services: mysql: image: mysql:8.0.33 container_name: mysql-cdc environment: MYSQL_ROOT_PASSWORD: root123 MYSQL_DATABASE: shop ports: - 3306:3306 command: - --binlog_formatROW - --server-id1 - --log-binmysql-bin kafka: image: bitnami/kafka:3.4 container_name: kafka-cdc ports: - 9092:9092 environment: KAFKA_CFG_NODE_ID: 0 KAFKA_CFG_PROCESS_ROLES: controller,broker KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0kafka:9093 KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT用 Docker 跑 MySQL 有一点要注意容器里的 MySQL 默认可能没开 binlog所以要在command或者 my.cnf 中显式指定binlog_formatROW。很多新手在这里翻车FlinkCDC 一直提示读取不到 binlog排查半天发现容器里的 MySQL 压根没开。2.3 最终环境清单整个链路的环境清单我列成一张表方便参照操作系统CentOS 7 / macOS 均可JDK1.8Flink 1.18 建议 JDK 11构建工具Maven 3.6开发语言Java下游 Topicods_table_changes分区数 62.4 引入 Maven 依赖时注意排除冲突pom.xml 里有一个细节flink-connector-mysql-cdc会传递引入 Debezium 相关依赖这些依赖和你 Flink 环境里的某些组件可能会冲突。我的做法是显式排除部分传递依赖只保留核心的 connector。dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-mysql-cdc/artifactId version2.4.2/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version1.18.0/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version1.18.1/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version1.18.1/version scopeprovided/scope /dependency如果你用 Flink 1.18 但引入了flink-connector-mysql-cdc:3.0.0源码里会用到一些新 API网上很多 2.x 的教程直接抄就会编译不过。所以版本一定要先定死别边写边升。3. MySQL 端先说清两件事binlog 配置和 CDC 账号权限3.1 修改 MySQL 配置开启 binlogMySQL 8.0 默认可能没开 binlog或者开的是 STATEMENT 格式。FlinkCDC 不能用。修改配置文件/etc/my.cnf在[mysqld]段下加[mysqld] server-id1 log-binmysql-bin binlog_formatROW binlog_row_imageFULL gtid_modeON enforce_gtid_consistencyONserver-id每个 MySQL 实例必须唯一尤其集群部署时不能重复。binlog_row_imageFULL让 binlog 记录每一行的完整前后镜像FlinkCDC 才能拿到 before 和 after。如果不设UPDATE 事件里可能只记录被修改的字段下游做对比分析时会缺数据。gtid_modeON建议开启。GTID 能让 FlinkCDC 在任务重启后更精准地定位位点尤其是多实例切换时。修改完重启 MySQLservice mysqld restart然后验证是否生效SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format; SHOW VARIABLES LIKE binlog_row_image;这三个值必须分别是ON、ROW、FULL缺一个都会导致后面的链路出问题。3.2 创建 CDC 专用账号并授权FlinkCDC 不建议直接用 root 连 MySQL一个是安全性问题另一个是权限太宽。官方需要的最小权限如下CREATE USER cdc% IDENTIFIED BY cdc123; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc%; FLUSH PRIVILEGES;SELECT做初始快照时读取表数据。RELOAD某些情况下需要刷新权限或状态。SHOW DATABASESFlinkCDC 枚举表结构时需要。REPLICATION SLAVE和REPLICATION CLIENT订阅 binlog 的核心权限。这里踩过一个坑一开始只给了SELECT结果 FlinkCDC 启动后一直报权限错误说无法注册为 slave。日志里没有明确提示缺哪个权限最后逐项排查才发现少了REPLICATION SLAVE。3.3 用真实数据验证 binlog 在产生可以在 MySQL 里执行一条 UPDATE然后查看 binlog 事件UPDATE shop.orders SET status paid WHERE id 100; SHOW BINLOG EVENTS IN mysql-bin.000001 LIMIT 10;如果能看到Update_rows事件说明 binlog 已经正常记录变更。我习惯在搭环境时先做这一步确认 binlog 没问题再继续写 Flink 代码不然问题堆在一起很难定位。4. Kafka 侧的准备工作Topic、写入方式和几个核心参数4.1 创建 Topic 时分区数怎么定Kafka 这边不需要太多花活先创建一个 topickafka-topics.sh --create \ --topic ods_table_changes \ --partitions 6 \ --replication-factor 1 \ --bootstrap-server localhost:9092分区数我建议按下游消费能力和数据量综合定。replication-factor如果是单机测试就设 1生产环境至少 3。有一个很关键的思路我一开始纠结到底是按表分多个 topic还是所有表写入同一个 topic。最后选了后者所有表的变更数据统一写到ods_table_changes在消息体里带上 database 和 table 字段下游自己过滤。原因是 Kafka 的 topic 数量一多消费端维护成本翻倍而这个场景对吞吐要求没那么夸张一个 topic 加足够分区完全够用。4.2 写入 Kafka 的两种 API别再用老的了Flink 社区现在推荐使用KafkaSink而不是老的FlinkKafkaProducer。FlinkKafkaProducer在 Flink 1.15 开始被标记为弃用虽然还能用但它的事务机制和新的KafkaSink相比有差距尤其是对 EXACTLY_ONCE 的支持细节上。KafkaSink用法更简洁KafkaSinkString kafkaSink KafkaSink.Stringbuilder() .setBootstrapServers(localhost:9092) .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic(ods_table_changes) .setValueSerializationSchema(new SimpleStringSchema()) .build()) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix(flink-cdc-txn) .setProperty(transaction.timeout.ms, 900000) .build();老 API 里setProducerConfig那一套在新的KafkaSink里统一用setProperty设置。如果有老项目迁移过来注意改成这种写法。4.3 三个必调的生产者参数DeliveryGuarantee.EXACTLY_ONCE和AT_LEAST_ONCE的选择取决于你的数据管道允不允许重复。DeliveryGuarantee.EXACTLY_ONCEKafka producer 会使用事务消息在 checkpoint 完成时才对外可见。这个模式最安全但对transaction.timeout.ms要求高。DeliveryGuarantee.AT_LEAST_ONCE没有事务开销吞吐更高但故障恢复时可能重复写入。我线上用的是 EXACTLY_ONCE并显式设置了transaction.timeout.ms为 15 分钟。默认值在很多版本里只有 1 分钟如果你 checkpoint 间隔设置较长或者某个批次积压了数据事务还没完成就被 broker 判定超时作业会一直报错。还有一个参数容易被忽略acksall。虽然KafkaSink在 EXACTLY_ONCE 模式下会自动处理但如果你在setProperty里覆盖了acks千万别写成acks1否则可能以牺牲一致性为代价换来吞吐。5. DataStream 方式的核心实现从 Source 到 Sink 的完整代码5.1 构建 MySQL CDC SourceMySqlSource是链路的数据入口。MySqlSourceString mySqlSource MySqlSource.Stringbuilder() .hostname(localhost) .port(3306) .username(cdc) .password(cdc123) .databaseList(shop) .tableList(shop.orders) .serverTimeZone(Asia/Shanghai) .deserializer(new OrderJsonDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build();databaseList和tableList是白名单只监听指定库表。tableList的格式是库名.表名必须带库名这个很多人会写错。startupOptions有几种选择StartupOptions.initial()先做一次全量快照再无缝切换到增量。生产环境最常用。StartupOptions.latest()只从当前时间往后解析历史数据不管。开发调试时方便。StartupOptions.earliest()从最早可用的 binlog 位置开始。适合重放场景但 MySQL 的 binlog 如果没有保留足够长时间会直接失败。开发阶段用latest()快正式上线我会用initial()确保 Kafka 里不会缺历史数据。5.2 自定义反序列化器不要偷懒用 String很多人图方便直接用StringDebeziumDeserializationSchema拿到的是整个 Debezium JSON包含一堆中间字段。虽然能跑但下游解析麻烦而且 with 状态里还夹着大量没用的元信息。我更建议自己写一个反序列化器只提取需要的数据。public class OrderJsonDeserializationSchema implements DebeziumDeserializationSchemaString { Override public void deserialize(SourceRecord record, CollectorString out) { Struct value (Struct) record.value(); if (value null) { return; } Struct source value.getStruct(source); String database source.getString(db); String table source.getString(table); String op value.getString(op); Struct after value.getStruct(after); String afterJson after null ? null : after.toString(); Struct before value.getStruct(before); String beforeJson before null ? null : before.toString(); JSONObject result new JSONObject(); result.put(database, database); result.put(table, table); result.put(op, op); result.put(before, beforeJson); result.put(after, afterJson); result.put(ts_ms, value.getInt64(ts_ms)); out.collect(result.toJSONString()); } Override public TypeInformationString getProducedType() { return BasicTypeInfo.STRING_TYPE_INFO; } }before和after都转成 JSON 字符串保留这样下游既能拿到变更后的数据也能对比变更前的值。op字段的取值通常是ccreate、uupdate、ddelete、r快照时的读取。5.3 组装完整作业并启动public class MysqlCdcToKafkaJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 开启 checkpointEXACTLY_ONCE 依赖它 env.enableCheckpointing(10000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5000); env.getCheckpointConfig().setCheckpointTimeout(60000); MySqlSourceString mySqlSource MySqlSource.Stringbuilder() .hostname(localhost) .port(3306) .username(cdc) .password(cdc123) .databaseList(shop) .tableList(shop.orders) .serverTimeZone(Asia/Shanghai) .deserializer(new OrderJsonDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build(); KafkaSinkString kafkaSink KafkaSink.Stringbuilder() .setBootstrapServers(localhost:9092) .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic(ods_table_changes) .setValueSerializationSchema(new SimpleStringSchema()) .build()) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix(flink-cdc-txn) .setProperty(transaction.timeout.ms, 900000) .build(); DataStreamSourceString stream env.addSource(mySqlSource); stream.sinkTo(kafkaSink); env.execute(mysql-cdc-to-kafka); } }这段代码可以直接跑前提是 MySQL、Kafka、Flink 版本和我上面列的一致。sinkTo是KafkaSink对应的方法老 API 才用addSink。5.4 为什么用 DataStream 而不是 SQL ClientFlink CDC 也支持纯 SQL 方式写 Kafka语法上几行就够CREATE TABLE orders_source (...); CREATE TABLE kafka_sink (...); INSERT INTO kafka_sink SELECT ... FROM orders_source;但 DataStream 方式在这几个方面更有优势反序列化逻辑可控。你能在写入 Kafka 前对消息做字段裁剪、格式转换、路由选择。多路输出方便。一个 CDC 源可以同时写到多个 Kafka topic 或别的 SinkSQL 里要写多个INSERT INTO语句逻辑分散。复杂状态处理时更顺手。比如要根据变更类型做不同的清洗规则Java 代码比 SQL 直观得多。对于中小团队DataStream 更符合 Java 工程师的日常开发习惯。SQL Client 适合快速验证做正式管道我建议用 DataStream。6. 跑通之后必须动的手脚checkpoint、并行度与时区6.1 checkpoint 是 EXACTLY_ONCE 的地基很多人以为设置了DeliveryGuarantee.EXACTLY_ONCE就万事大吉其实不是。Flink 的 KafkaSink 在 EXACTLY_ONCE 模式下会把消息写入 Kafka 事务但这个事务只有在 checkpoint 成功时才 commit。换句话说如果不开启 checkpoint或者 checkpoint 一直失败Kafka 里的数据要么不出现要么出现得很慢。我线上的 checkpoint 配置env.enableCheckpointing(10000); env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5000); env.getCheckpointConfig().setCheckpointTimeout(60000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);MinPauseBetweenCheckpoints和setMaxConcurrentCheckpoints(1)配套使用避免频繁 checkpoint 给 MySQL 和 Kafka 造成额外压力。6.2 Source 并行度的问题Flink CDC 2.x 的MySqlSource是普通SourceFunction并行度只能为 1。这是设计限制因为读取一个 MySQL 实例的 binlog 本质上是单连接、严格有序的。你在addSource后面调setParallelism(2)是没有意义的甚至可能报错。那吞吐不够怎么办两个思路按表拆分多个 CDC 作业各自订阅不同的表均衡分配给多个 TaskManager。如果表实在太大、变更太频繁可以考虑 Flink CDC 3.0 的 YAML pipeline 方式它支持并行读取的部分场景但需要重新评估 API 和依赖兼容。我在生产环境用了方案一一个 Flink 作业订阅核心大表另一个作业订阅其余小表。这样单个作业的压力可控出问题也能快速隔离。6.3 时区问题比想象的隐蔽MySQL 的TIMESTAMP类型在存储时会转成 UTC读取后再按会话时区转回本地时间。FlinkCDC 在构建 Source 时虽然可以指定serverTimeZone(Asia/Shanghai)但 Debezium 内部处理 binlog 时很多时间字段会转成 epoch 毫秒或者 ISO8601 字符串。如果下游直接将这些值展示成当地时间很可能会比实际慢 8 个小时。处理方式有两种在自定义反序列化器里遇到ts_ms或者 binlog 里的时间字段显式用Instant.ofEpochMilli(x).atZone(ZoneId.of(Asia/Shanghai))格式化。下游消费 Kafka 时统一按照 UTC 解析展示端再转时区。我建议选第一种让 Kafka 里的消息永远是“北京时间”的可读字符串省得下游每个消费端各自转换容易转错。6.4 低活动表与心跳机制如果同步的某张表一整天都没数据变更binlog 上长时间没有事件FlinkCDC 的连接可能处于一种“几乎静止”的状态。此时如果 MySQL 端发生网络切换或连接超时Flink 可能在较长时间后才感知到。Flik CDC 2.4 里提供了heartbeatInterval配置建议设置一个合理值比如 30 秒.heartbeatInterval(Duration.ofSeconds(30))这会让 FlinkCDC 每隔一段时间在 binlog 里产生一个心跳事件既是探测连接是否存活的信号也能避免长时间无数据时低延迟指标完全失真。这个参数最初我没加结果有一次 MySQL 主从切换FlinkCDC 过了好几分钟才报错重启心跳加上之后故障发现时间缩短到 30 秒内。7. 上线后踩过的坑和排查链路7.1 server-id 冲突多个作业连同一个 MySQL 时必须避开这是我踩得最惨的一次。项目里两个 FlinkCDC 作业同时连同一个 MySQL 实例一个同步订单表一个同步商品表。第一个作业跑了一个月都正常第二个作业上线后第一个作业开始频繁断开日志里出现Slave can not handle replication events with the checksum that master is configured to log或者连接被踢掉的错误。原因就是两个作业复用了同一个server-id。MySQL 认为这是同一个复制连接后建立的连接会把先建立的挤掉。解决方法是每个作业分配不同的serverId区间// 作业A .serverId(5400-5404) // 作业B .serverId(5500-5504)不只是 Flink 作业之间如果还有其他 Canal 或 Debezium 实例也在这个 MySQL 上拉 binlog一样要在配置里错开。这个冲突不会在你开发时暴露往往是在第二个任务加上线后才出现排查起来特别容易被忽略。7.2 DDL 变更把下游 JSON 解析搞挂了上线后某一天下游突然报大量 JSON 解析异常。查到最后是我在订单表上加了一个字段。FlinkCDC 只负责把 binlog 里的变更数据发出去它不会主动通知你“表结构变了请修改反序列化器”。如果你在自定义反序列化器里写死了字段比如每次都取value.getStruct(after).getString(order_name)新增字段后after结构变了解析就可能直接抛异常。我当时排查的链路是下游消费程序日志先出现org.apache.kafka.common.errors.SerializationException接着 Kafka 里看到正常消息和异常消息混杂。最后定位到 CDC 的反序列化器代码发现它写死了一个枚举字段来路由消息。处理方式有两个反序列化器里尽量把after整体作为一个 JSON 字符串保留不要提前拆开下游拿到后再按需解析。这样上游表结构新增字段不会影响 CDC 作业。如果一定要在 CDC 里做字段映射用JSONObject动态解析取字段前先containsKey判断。7.3 checkpoint 超时与 Kafka 事务超时有一次压测单条 SQL 批量更新了十几万行数据。FlinkCDC 把这批变更全部读出来写入 Kafka 时用了事务。结果 checkpoint 一直失败报错信息指向 Kafka 事务超时。排查链路如下打开 Flink UI看到Checkpoint页面全是EXPIRED。打开异常日志找到Transaction timed out due to timeout。确认transaction.timeout.ms默认值太小。调大 producer 的transaction.timeout.ms同时确保它小于 broker 的transaction.max.timeout.ms默认 15 分钟。最终我设置的是900000毫秒也就是 15 分钟。这个值不是瞎调的它必须大于一次 checkpoint 可能经历的最大间隔。如果一张大表的 snapshot 阶段本身耗时很长还要结合checkpoint.timeout和minPauseBetweenCheckpoints综合评估。7.4 Kafka 控制台消费命令为什么“启动一次会一直运行”很多初接触 Kafka 的同事问过我一个问题执行kafka-console-consumer.sh后窗口就停在那里不动了是不是命令卡死了这不是卡死是正常现象。控制台消费者默认会持续监听 topic 的新消息没有消息时就一直等待如果要退出按CtrlC就行。kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic ods_table_changes \ --from-beginning在调试 CDC 链路时我会开一个这样的窗口挂在旁边然后在 MySQL 里执行一条 UPDATE窗口里马上能看到新消息。这个小技巧排查问题非常高效。顺便提醒如果加了--from-beginning会从 topic 最早的消息开始打印历史消息会刷屏只想看新增消息去掉这个参数即可。7.5 Kafka 消息延迟高先看 lag 再看生产者参数同步链路刚上线时发现 Kafka 里的数据比 MySQL 实际变更慢了十几秒。第一反应是 Flink 处理慢打开 UI 一看全程没有反压。后来查了 consumer group 的 lag发现消息明明在 topic 里但下游消费滞后严重问题出在下游。这也说明排查顺序很重要先定位瓶颈在哪一段再动手改配置。如果确认瓶颈在 Flink 写入 Kafka 这一段优先检查这几个配置linger.msKafkaSink默认是 0消息立即发送延迟最低但吞吐一般。如果调成 10ms会积攒一小批再发吞吐提高但延迟增加。batch.size适当调大可以减少网络请求次数。分区数分区太少单分区写入压力大也会造成延迟。我最终用的方案是linger.ms5、batch.size64KB实测在延迟和吞吐之间比较平衡。7.6 大事务会让整条链路出现尖刺有一次业务团队手动跑了一条 UPDATE一次性改了全表 80 万行。FlinkCDC 需要把这 80 万行的变更全部读出来Kafka 侧短时间内写入量暴增下游消费 lag 一下飙到几十万。这个不是故障但确实会“吓人一跳”。如果你有这样的重操作场景建议在上游业务层面把大事务拆成小批量比如每 5000 行提交一次。否则无论 Flink 还是 Kafka 都要为这种极端流量预留 buffer资源浪费是长期的。结尾一点实战体会这套 FlinkCDC Kafka 的链路我前后跑了大半年最大的体会是先把环境验证清楚再写业务代码。很多人一上来就写 Flink 作业结果 binlog 没开、server-id 冲突、Kafka 事务超时这些基础问题全堆在一起反而找不准方向。我的习惯是先确认 MySQL 的 binlog 能正常产生再用 Kafka 控制台消费者确认 topic 能收消息最后才写 DataStream 代码。每一步都用最小方式验证过再往下走出问题的时候能省下大量排查时间。另外上线前一定要把 checkpoint 监控加上看 checkpoint 的成功率和耗时这比什么指标都更能反映整条链路的健康度。