You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink 1.15.1使用Kafka CDC表连接器重启后Sink出现重复消息问题

问题复现条件

  • 作业逻辑:从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执行直接报错,需要先修复。

修复方案

  1. 调整检查点与Savepoint配置
    将检查点留存策略修改为取消作业时保留,避免状态数据被自动清理:
    configuration.setString("execution.checkpointing.externalized-checkpoint-retention", "RETAIN_ON_CANCELLATION");
    // 单独配置Savepoint存储路径,便于运维管理
    configuration.setString("state.savepoints.dir", "file:///tmp/savepoints/");
    
    取消作业时主动触发Savepoint,重启作业必须通过-s参数指定本次触发的Savepoint路径,或使用最新留存的检查点路径启动,禁止无状态启动作业。
  2. 为所有Kafka连接器配置恰好一次语义
    在三个Upsert Kafka表的WITH参数中添加如下配置,开启事务支持实现端到端恰好一次:
    'delivery.guarantee' = 'exactly-once',
    'sink.transaction.timeout.ms' = '900000'
    

    注意:事务超时时间需要大于检查点间隔(当前配置为3min),同时小于Kafka集群配置的transaction.max.timeout.ms(默认值为15min),避免事务超时失效。

  3. 显式配置源表消费起始位点
    两个源表添加启动模式配置,避免无状态启动时的非预期消费:
    'scan.startup.mode' = 'earliest-offset'
    
    该配置仅在作业无状态首次启动时生效,后续从状态恢复时会自动使用状态中存储的消费位点,不受该配置影响。
  4. 修复DDL语法错误:删除test2表定义中correlationId STRING行末尾多余的逗号。

验证方法

  • 修正代码与配置后重新部署作业,等待作业完成第一次检查点,或主动触发一次Savepoint
  • 停止作业,使用刚才生成的Savepoint/检查点路径指定恢复启动
  • 校验下游Kafka主题消息量,不会再出现翻倍重复问题,相同主键的消息仅保留最新版本,符合Upsert语义预期。

内容的提问来源于stack exchange,提问作者user3332207

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 04:57:24