Flink 1.15.1使用Kafka CDC表连接器重启后Sink出现重复消息问题
Flink 1.15.1 作业重启后下游Kafka出现重复消息排查
问题复现条件
- 作业逻辑:从2个Kafka主题读取CDC消息,完成数据转换后写入下游Kafka主题,配置状态后端与检查点,计划通过Savepoint实现作业取消后的状态恢复
- 异常表现:作业首次运行下游目标主题共写入10条消息,重启作业后主题消息总量增长至20条,出现全量重复
根因定位
- 检查点留存配置错误:当前配置
execution.checkpointing.externalized-checkpoint-retention为DELETE_ON_CANCELLATION,作业主动取消时会自动删除所有已生成的检查点。如果取消作业时未主动触发Savepoint,或重启时未指定有效Savepoint/检查点路径,作业会以无状态方式启动,按照配置的消费起始位点重新拉取数据,导致重复写入。 - Upsert Kafka连接器投递语义配置缺失:所有Kafka源表、Sink表均未显式配置投递保障级别,Flink 1.15版本Upsert Kafka连接器默认投递保障为
at-least-once,即使全局检查点配置为EXACTLY_ONCE模式,Sink未开启两阶段提交的事务支持,重启场景下依然会产生重复消息。 - 源表消费起始位点未显式配置:两个Kafka源表均未指定
scan.startup.mode,无状态启动时如果消费组位点不存在,会按照默认规则拉取历史数据,进一步放大重复问题。 - 代码存在语法错误:test2表DDL中
correlationId STRING字段后多余一个逗号,会导致SQL执行直接报错,需要先修复。
修复方案
- 调整检查点与Savepoint配置
将检查点留存策略修改为取消作业时保留,避免状态数据被自动清理:
取消作业时主动触发Savepoint,重启作业必须通过configuration.setString("execution.checkpointing.externalized-checkpoint-retention", "RETAIN_ON_CANCELLATION"); // 单独配置Savepoint存储路径,便于运维管理 configuration.setString("state.savepoints.dir", "file:///tmp/savepoints/");-s参数指定本次触发的Savepoint路径,或使用最新留存的检查点路径启动,禁止无状态启动作业。 - 为所有Kafka连接器配置恰好一次语义
在三个Upsert Kafka表的WITH参数中添加如下配置,开启事务支持实现端到端恰好一次:'delivery.guarantee' = 'exactly-once', 'sink.transaction.timeout.ms' = '900000'注意:事务超时时间需要大于检查点间隔(当前配置为3min),同时小于Kafka集群配置的
transaction.max.timeout.ms(默认值为15min),避免事务超时失效。 - 显式配置源表消费起始位点
两个源表添加启动模式配置,避免无状态启动时的非预期消费:
该配置仅在作业无状态首次启动时生效,后续从状态恢复时会自动使用状态中存储的消费位点,不受该配置影响。'scan.startup.mode' = 'earliest-offset' - 修复DDL语法错误:删除test2表定义中
correlationId STRING行末尾多余的逗号。
验证方法
- 修正代码与配置后重新部署作业,等待作业完成第一次检查点,或主动触发一次Savepoint
- 停止作业,使用刚才生成的Savepoint/检查点路径指定恢复启动
- 校验下游Kafka主题消息量,不会再出现翻倍重复问题,相同主键的消息仅保留最新版本,符合Upsert语义预期。
内容的提问来源于stack exchange,提问作者user3332207
相关产品推荐
相关产品推荐

