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

Kafka Streams分区消息阻塞问题排查求助(事务未提交)

Kafka Streams事务消息阻塞根因分析与解决方案

问题现象

  • 基于Kafka Streams/Spring Boot/Spring Kafka的应用出现异常:某主题单个分区内数千条消息阻塞
  • read_committed隔离级别的消费者无法读取这些消息,仅read_uncommitted消费者可读取,推测消息已生产但事务未提交

环境版本

  • kafka-clients、kafka-streams:3.2.3
  • Spring Boot:2.7.9
  • Spring Kafka:2.9.1
  • Kafka Broker:3.4.0

核心配置

poll.ms = 10
commit.interval.ms = 0ms
topology.optimization = all
processing.guarantee = exactly_once_v2
num.standby.replicas = 1
metrics.recording.level = INFO
num.stream.threads = 64
replication.factor = 1
acks = all
linger.ms = 0
retries = 2147483647
max.block.ms = 60000
metadata.max.age.ms = 300000
max.poll.interval.ms = 60000
max.poll.records = 30
session.timeout.ms = 10000
fetch.max.wait.ms = 500
fetch.min.bytes = 1
heartbeat.interval.ms = 3000
allow.auto.create.topics = false
auto.offset.reset = latest
max.task.idle.ms = 10

异常触发场景与日志分析

异常发生时,生产者所在Kubernetes节点正被驱逐终止,Pod内捕获到如下关键异常:

网络异常阶段

2023-08-22 01:13:34.026  WARN [gateway-united-kafka-service,,] 1 --- [ead-64-producer] o.a.k.clients.producer.internals.Sender  : [Producer clientId=wagers-united-transactions-stream-c064ad3f-ba6a-4371-9e21-33e14af4a87a-StreamThread-64-producer, transactionalId=wagers-united-transactions-stream-502b91fe-fe5e-4af1-be54-bd1b25181882-64] Got error produce response with correlation id 10156 on topic-partition wagers-united-transactions-recovery-6, retrying (2147480373 attempts left). Error: NETWORK_EXCEPTION. Error Message: Disconnected from node 3

2023-08-22 01:13:34.026  WARN [gateway-united-kafka-service,,] 1 --- [ead-64-producer] o.a.k.clients.producer.internals.Sender  : [Producer clientId=wagers-united-transactions-stream-c064ad3f-ba6a-4371-9e21-33e14af4a87a-StreamThread-64-producer, transactionalId=wagers-united-transactions-stream-502b91fe-fe5e-4af1-be54-bd1b25181882-64] Received invalid metadata error in produce request on partition wagers-united-transactions-recovery-6 due to org.apache.kafka.common.errors.NetworkException: Disconnected from node 3. Going to request metadata update now

2023-08-22 01:14:04.142  WARN [gateway-united-kafka-service,,] 1 --- [ead-64-producer] o.a.k.clients.producer.internals.Sender  : [Producer clientId=wagers-united-transactions-stream-c064ad3f-ba6a-4371-9e21-33e14af4a87a-StreamThread-64-producer, transactionalId=wagers-united-transactions-stream-502b91fe-fe5e-4af1-be54-bd1b25181882-64] Got error produce response with correlation id 10159 on topic-partition wagers-united-transactions-recovery-6, retrying (2147480372 attempts left). Error: NETWORK_EXCEPTION. Error Message: Disconnected from node 3

2023-08-22 01:14:04.142  WARN [gateway-united-kafka-service,,] 1 --- [ead-64-producer] o.a.k.clients.producer.internals.Sender  : [Producer clientId=wagers-united-transactions-stream-c064ad3f-ba6a-4371-9e21-33e14af4a87a-StreamThread-64-producer, transactionalId=wagers-united-transactions-stream-502b91fe-fe5e-4af1-be54-bd1b25181882-64] Received invalid metadata error in produce request on partition wagers-united-transactions-recovery-6 due to org.apache.kafka.common.errors.NetworkException: Disconnected from node 3. Going to request metadata update now
  • 节点驱逐导致网络断开,生产者持续重试发送消息,但无法连接到Broker节点3

生产者围栏与事务异常阶段

2023-08-22 01:14:18.614 ERROR [gateway-united-kafka-service,64e2cdf1f01ed187e85e2a18a3d3615e,2f422c46d406a8cd] 1 --- [read-1-producer] o.a.k.s.p.internals.RecordCollectorImpl  : stream-thread [wagers-united-transactions-recovery-stream-gateway-united-kafka-66db75bbbf-ccgzv-StreamThread-1] task [0_4] Error encountered sending record to topic wagers-united-transactions-done for task 0_4 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

org.apache.kafka.common.errors.InvalidProducerEpochException: Producer attempted to produce with an old epoch.

2023-08-22 01:14:18.614 ERROR [gateway-united-kafka-service,64e2cdf1f01ed187e85e2a18a3d3615e,e44a9f539cb474cb] 1 --- [read-2-producer] o.a.k.s.p.internals.RecordCollectorImpl  : stream-thread [wagers-united-transactions-recovery-stream-gateway-united-kafka-66db75bbbf-ccgzv-StreamThread-2] task [0_4] Error encountered sending record to topic wagers-united-transactions-done for task 0_4 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

org.apache.kafka.common.errors.InvalidProducerEpochException: Producer attempted to produce with an old epoch.

2023-08-22 01:14:18.624  INFO [gateway-united-kafka-service,,] 1 --- [ead-64-producer] o.a.k.c.p.internals.TransactionManager   : [Producer clientId=wagers-united-transactions-stream-c064ad3f-ba6a-4371-9e21-33e14af4a87a-StreamThread-64-producer, transactionalId=wagers-united-transactions-stream-502b91fe-fe5e-4af1-be54-bd1b25181882-64] 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.
  • 新Pod启动后,Kafka Streams重新分配任务,新生产者使用相同transactionalId注册,触发**生产者围栏(Producer Fenced)**机制,旧生产者被标记为失效
  • 旧生产者在网络恢复前已发送的消息处于未提交状态,且旧生产者被围栏后无法完成事务提交/回滚;而新生产者无法感知这些未提交事务,导致消息长期阻塞

根因总结

  1. Kubernetes节点驱逐导致生产者与Broker网络中断,消息发送重试但无法完成
  2. **Exactly-Once语义(exactly_once_v2)**依赖事务机制,旧生产者的事务在网络中断时处于悬停状态
  3. 生产者围栏机制触发后,旧生产者无法继续操作事务,新生产者不会处理旧事务的遗留消息
  4. 主题副本配置缺陷:replication.factor=1,无副本冗余,Broker节点故障或网络中断时无备用节点承接,加剧事务悬停风险

解决方案与预防措施

短期修复

  • 使用kafka-consumer-groups工具或自定义read_uncommitted消费者将阻塞消息迁移至新主题,恢复业务流程
  • 手动清理悬停事务:通过Kafka Broker的TransactionCoordinator查询未完成事务,强制提交或回滚(需谨慎操作,避免数据不一致)

长期预防

  1. 调整副本配置:将主题replication.factor设置为3(生产环境标准配置),提升Broker侧冗余能力,避免单节点故障导致网络中断
  2. 优化Kubernetes驱逐策略:配置Pod的terminationGracePeriodSeconds(建议30s以上),给生产者足够时间完成事务提交/回滚;添加preStop钩子,调用Kafka Streams的close()方法,确保优雅关闭
  3. 调整生产者重试与事务配置:
    • 降低retries值(无需设置为最大值),避免长期无意义重试占用资源
    • 设置transaction.timeout.ms(默认15min),确保悬停事务能被自动清理
  4. 升级Kafka版本:3.2.3存在部分事务相关的Bug(如悬停事务清理不及时),建议升级至3.4.x及以上版本,与Broker版本保持一致
  5. 监控事务状态:添加Kafka Streams事务指标监控(如transaction-aborted-total、transaction-committed-total),及时发现悬停事务

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 13:50:53