Flink Kafka EXACTLY_ONCE报InvalidProducerEpochException异常
Flink 1.15集成Kafka抛出InvalidProducerEpochException异常排查
问题背景
- 基于Flink 1.15.0版本完成Kafka Source、Sink集成开发,Sink配置投递保证级别为EXACTLY_ONCE
- 当前Kafka生产者配置
transaction.timeout.ms=60000(1分钟),应用运行一段时间后意外终止,抛出异常:
Caused by: org.apache.kafka.common.errors.InvalidProducerEpochException: Producer attempted to produce with an old epoch.
- 已知可将
transaction.timeout.ms调大到最大值900000ms(15分钟),但判断该调整无法彻底解决问题。
相关实现代码
KafkaSource实现
public static KafkaSource<String> kafkaSource(String bootstrapServers, String topic, String groupId) { return KafkaSource.<String>builder() .setBootstrapServers(bootstrapServers) .setTopics(topic) .setGroupId(groupId) .setStartingOffsets(OffsetsInitializer.latest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); }
KafkaSink实现
public static KafkaSink<WordCountPojo> kafkaSink(String brokers, String topic, Properties producerProperties) { return KafkaSink.<WordCountPojo>builder() .setKafkaProducerConfig(producerProperties) .setBootstrapServers(brokers) .setRecordSerializer(new WordCountPojoSerializer(topic)) .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix(UUID.randomUUID().toString()) .build(); }
自定义Kafka序列化器实现
public class WordCountPojoSerializer implements KafkaRecordSerializationSchema<WordCountPojo> { private String topic; private ObjectMapper mapper; public WordCountPojoSerializer(String topic) { this.topic = topic; this.mapper = new ObjectMapper(); } @Override public ProducerRecord<byte[], byte[]> serialize(WordCountPojo wordCountPojo, KafkaSinkContext kafkaSinkContext, Long timestamp) { try { byte[] serializedValue = mapper.writeValueAsBytes(wordCountPojo); return new ProducerRecord<>(topic, null, timestamp, null, serializedValue); } catch (JsonProcessingException e) { return null; } } }
异常根因
你之前怀疑的「两个事务间隙生产消息导致异常」判断不成立,核心触发原因有两个,其中第二个是致命配置错误:
- 事务超时与Checkpoint周期不匹配:Kafka Broker会主动中止超过
transaction.timeout.ms时长未提交的事务,同时将对应生产者的epoch值加1。如果Flink作业的Checkpoint间隔+Checkpoint最大执行耗时超过你设置的1分钟,事务会被Broker强制中止,后续Flink拿着旧epoch尝试提交事务、生产消息时就会抛出该异常。 - 事务ID前缀配置错误:Flink Kafka Sink实现EXACTLY_ONCE语义的核心前提是使用固定的事务ID前缀:作业Failover、Task重启时,会基于「固定前缀+子任务编号」生成和之前一致的事务ID,主动回滚之前未完成的残留事务。你当前用
UUID.randomUUID().toString()生成随机前缀,每次重启后前缀完全变化,新的生产者实例无法感知、回滚之前的残留事务,旧事务被Broker超时中止后,持有过期epoch的生产者发起请求就会触发异常。 - 额外代码隐患:自定义序列化器在JSON序列化失败时直接返回
null,会导致数据静默丢弃,违反EXACTLY_ONCE语义要求。
解决方案
按优先级调整配置即可彻底解决问题,仅调大transaction.timeout.ms无法解决核心问题:
- 首先修正事务ID前缀配置:删除随机UUID生成逻辑,为每个作业配置一个全局唯一、固定不变的字符串前缀,例如
flink-wordcount-stat-sink,保证作业重启前后前缀一致。 - 对齐事务超时与Checkpoint配置:
- 按
Checkpoint间隔 + 最大Checkpoint耗时 + 30~60s冗余设置transaction.timeout.ms值,例如Checkpoint间隔设为120s、最大Checkpoint耗时预期30s,就将该值设为180000(3分钟),注意该值必须小于Kafka Broker端transaction.max.timeout.ms配置(默认900000ms即15分钟,不建议随意调大Broker端该配置,会导致事务元数据长期占用Broker内存)。 - 配置Flink Checkpoint超时时间,保证Checkpoint超时阈值小于设置的
transaction.timeout.ms,避免Checkpoint长时间阻塞导致事务超时。
- 按
- 修复序列化器逻辑:JSON序列化失败时不要返回null,直接抛出RuntimeException触发作业失败告警,或通过侧输出流收集脏数据单独处理,避免数据静默丢失。
- 可选优化:如果作业对延迟不敏感,将Checkpoint最大并发数设为1,避免多并发Checkpoint阻塞事务提交流程;同时配置合理的重启策略,保证Failover时事务恢复逻辑正常执行。
内容的提问来源于stack exchange,提问作者Vin
相关产品推荐
相关产品推荐

