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

Kafka Streams回滚后报ProducerFencedException等错误的排查求助

Kafka Streams回滚后出现ProducerFenced/InvalidProducerEpoch异常的原因及恢复原理分析

错误场景回顾

基于Spring Kafka的Kafka Streams应用,部署存在Bug的版本后15-30分钟执行回滚,回滚后应用抛出以下异常:

  • ProducerFencedException:
[Producer clientId=myapplication-793b3cc4-9baf-4759-8516-649955fad000-StreamThread-2-producer, transactionalId=myapplication-793b3cc4-9baf-4759-8516-649955fad000-2] Transiting to fatal error state due to org.apache.kafka.common.errors.ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one. 
  • TaskMigratedException伴随InvalidProducerEpochException:
org.apache.kafka.streams.errors.TaskMigratedException: Error encountered sending record to topic myapplication-caliperEvents-changelog for task 0_9 due to:
org.apache.kafka.common.errors.InvalidProducerEpochException: Producer attempted to produce with an old epoch.
Written offsets would not be recorded and no more records would be sent since the producer is fenced, indicating the task may be migrated out; it means all tasks belonging to this thread should be migrated.
    at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.recordSendError(RecordCollectorImpl.java:215)
    at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.lambda$send$0(RecordCollectorImpl.java:196)
    at brave.kafka.clients.TracingCallback$DelegateAndFinishSpan.onCompletion(TracingCallback.java:58)
    at org.apache.kafka.clients.producer.KafkaProducer$InterceptorCallback.onCompletion(KafkaProducer.java:1350)
    at org.apache.kafka.clients.producer.internals.ProducerBatch.completeFutureAndFireCallbacks(ProducerBatch.java:273)
    at org.apache.kafka.clients.producer.internals.ProducerBatch.done(ProducerBatch.java:234)
    at org.apache.kafka.clients.producer.internals.ProducerBatch.completeExceptionally(ProducerBatch.java:198)
    at org.apache.kafka.clients.producer.internals.Sender.failBatch(Sender.java:758)
    at org.apache.kafka.clients.producer.internals.Sender.failBatch(Sender.java:743)
    at org.apache.kafka.clients.producer.internals.Sender.failBatch(Sender.java:695)
    at org.apache.kafka.clients.producer.internals.Sender.completeBatch(Sender.java:634)
    at org.apache.kafka.clients.producer.internals.Sender.lambda$null$1(Sender.java:575)
    at java.base/java.util.ArrayList.forEach(Unknown Source)
    at org.apache.kafka.clients.producer.internals.Sender.lambda$handleProduceResponse$2(Sender.java:562)
    at java.base/java.lang.Iterable.forEach(Unknown Source)
    at org.apache.kafka.clients.producer.internals.Sender.handleProduceResponse(Sender.java:562)
    at org.apache.kafka.clients.producer.internals.Sender.lambda$sendProduceRequest$5(Sender.java:836)
    at org.apache.kafka.clients.ClientResponse.onComplete(ClientResponse.java:109)
    at org.apache.kafka.clients.NetworkClient.completeResponses(NetworkClient.java:574)
    at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:566)
    at org.apache.kafka.clients.producer.internals.Sender.runOnce(Sender.java:328)
    at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:243)
    at java.base/java.lang.Thread.run(Unknown Source)
  • InvalidProducerEpochException:
[Producer clientId=myapplication-793b3cc4-9baf-4759-8516-649955fad000-StreamThread-2-producer, transactionalId=myapplication-793b3cc4-9baf-4759-8516-649955fad000-2] Transiting to abortable error state due to org.apache.kafka.common.errors.InvalidProducerEpochException: Producer attempted to produce with an old epoch. 

异常导致Persistent KV存储的记录无法删除,punctuator反复读取处理事件。临时解决方案为关闭Exactly-Once处理,运行约8小时后重新开启,应用恢复正常。


错误根源分析

这些异常均与Kafka Streams的**Exactly-Once语义(EOS)**绑定的事务机制直接相关:

  1. transactionalId与Producer Epoch冲突
    Kafka Streams的每个StreamThread会使用固定格式的transactionalId:<application.id>-<thread-number>,用于保证事务的幂等性和EOS语义。当部署新版本应用时,新的StreamThread会向Kafka Broker注册相同的transactionalId,并获取一个更高的Producer Epoch(Broker用来标识生产者实例版本的递增计数器)。
    回滚旧版本应用后,旧实例的StreamThread会尝试用之前的旧Epoch继续提交事务,Broker会判定该生产者为"过期实例",直接触发fence机制拒绝请求,抛出ProducerFencedException和InvalidProducerEpochException。

  2. 任务迁移失败导致状态存储更新阻塞
    当生产者被fence后,Kafka Streams会触发任务迁移逻辑,认为该任务已被其他实例接管。但由于新旧实例的transactionalId完全重复,任务迁移无法正常完成,导致状态存储的更新操作(如删除记录)无法提交到对应的changelog topic,最终KV存储中的记录持续存在,punctuator会反复读取处理。


临时恢复方案的原理

关闭Exactly-Once语义的作用

将processing.guarantee配置从exactly_once_v2改为at_least_once后:

  • Kafka Streams将不再启用事务机制,生产者不再绑定transactionalId和Producer Epoch,Broker不再进行epoch校验,应用可以正常写入changelog topic,完成KV存储的记录删除操作,终止事件的重复处理。

等待8小时后重新开启EOS的原因

Kafka Broker会保存每个transactionalId对应的最新Epoch和事务状态,这些元数据存储在内部的事务状态日志中。当某个transactionalId长时间没有活动(超过Broker配置的事务元数据保留时长,结合事务超时时间的叠加影响),对应的元数据会被Broker自动清理。
等待8小时是为了让Broker彻底清理掉之前新版本应用留下的transactionalId元数据,此时重新开启EOS,旧版本应用的StreamThread可以重新向Broker注册transactionalId并获取新的Epoch,不会再出现冲突,应用恢复正常运行。


总结

  • 核心冲突:回滚后的应用实例与之前的新版本实例使用相同的transactionalId,且Broker已记录更高的Producer Epoch,导致旧实例被Broker fence。
  • 恢复逻辑:关闭EOS绕过事务校验,解决状态存储更新阻塞;等待足够时间让Broker清理旧事务元数据,重新开启EOS时能正常注册新的事务上下文。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 11:47:05