前言

  最近在做反欺诈项目实时数据流搭建,需求是通过 Flink MySQL-CDC 抓取交易表增删改数据,同步投递到 Kafka 供下游风控引擎消费。本以为复制模板就能一键跑通,结果接连踩了三处典型大坑:initial全量模式跑完存量不切实时print()直接阻塞流式任务数据投递 Kafka 格式不对无法识别。把排查过程、根源、最终修复方案完整记录下来,给同样做 CDC 同步的同学避坑。

环境版本

  • Flink:1.17.0
  • MySQL CDC Connector:2.4.0
  • MySQL 8.0(开启 binlog,格式 ROW)
  • Kafka:3.2.x
  • 开发模式:Flink Table API Java

业务需求

同步Anti-Fraud库下user_trans_log交易表,全量同步历史交易数据,同步完成后持续监听 MySQL 实时新增 / 修改 / 删除操作,统一以 Debezium 标准 JSON 格式发送至 Kafka Anti_Fraud主题。

坑一:scan.startup.mode = initial 全量同步后不实时,latest-offset 却秒级监听

现象

  1. 配置latest-offset程序启动直接读取 binlog 末尾,MySQL 插入数据立刻推送到 Kafka,实时完全正常;
  2. 切换initial模式:程序只会疯狂拉取表里存量全量数据发送 Kafka,整张表数据没全部推送完毕前,新增的测试数据完全不会同步,看起来 “只能全量、不能实时”

问题根源

MySQL-CDC 的 initial 模式分为两个严格串行阶段

  1. 快照阶段:开启可重复读事务,JDBC 分页查询整张表所有存量数据,逐条输出;这个阶段 CDC 完全不会消费 binlog 事件,新写入 MySQL 的数据只会存进 binlog 积压。
  2. 增量 binlog 阶段必须整张表快照 100% 读取完成 + 成功生成 Checkpoint 位点,CDC 才会自动切换到 binlog 实时监听模式,开始消费积压 + 后续新的变更。

我一开始没开启 Checkpoint,大表快照读取缓慢,没有位点标记,导致快照结束后程序迟迟无法切换增量流,肉眼看就是 “全量跑完不动了”。

解决方案

解决方案

环境第一行强制开启 Checkpoint,给快照切换增量提供位点支撑

// 5秒一次检查点,保障快照完成后落地binlog位点
env.enableCheckpointing(5000);

坑二:一行 print 代码直接阻塞整个流式监听进程

tableEnv.executeSql(createCdcSql);
// 致命阻塞代码
tableEnv.executeSql("select * from user_trans_log_cdc ").print();
tableEnv.executeSql(createKafkaSinkSql);
tableEnv.executeSql("INSERT INTO trans_kafka_sink SELECT * FROM user_trans_log_cdc ");

现象

程序启动后 CDC 快照读取速度极慢,甚至卡死,Kafka 收不到数据,MySQL 新增数据完全无响应。

问题根源

  1. SELECT * FROM 流式CDC表本身是无限无界流,.print()是适配批处理的阻塞 API,会持续等待数据流 “结束”,但流式数据永远不会结束;
  2. 这条查询会单独启动一套 CDC 读取任务,和后面写入 Kafka 的 Insert 任务抢占 MySQL binlog 连接、数据库 IO 资源,双任务互相拖慢,快照阶段直接阻塞卡死,根本走不到实时监听环节。

正确写法

  1. 彻底删除单独的 select+print 语句,Flink Table API 流式任务只允许一条 Insert 写入逻辑作为程序主执行流;
  2. 调试查看数据不要用 print (),改用官方 print 连接器打印输出,无阻塞

坑三:CDC 原始数据直接丢 Kafka 无法解析,必须指定 format = debezium-json

踩坑经过

最开始偷懒没配置 format,默认用普通 json 格式投递 Kafka,下游消费拿到的只有 after 里的行数据,丢失操作类型(新增 / 修改 / 删除)、变更前数据等核心 CDC 元信息,无法区分 INSERT/UPDATE/DELETE。

关键说明

  1. debezium-json是 MySQL-CDC 配套专属序列化格式
    • 自带op操作标识:c = 新增、u = 更新、d = 删除、r = 全量快照数据
    • 自带before变更前数据、after变更后完整行
    • 兼容 Debezium 标准协议,所有大数据下游(Flink、Spark、ClickHouse、数仓)通用解析逻辑
  2. 如果只写普通'format' = 'json',只会单纯序列化 after 行对象,丢失 CDC 变更语义,失去同步删改的能力。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐