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

Kafka NotLeaderOrFollowerException仅在单个K8s集群出现的排查求助

问题描述

概述

现有两个K8s集群、一个3节点Broker的Kafka集群,以及一个6分区的Topic。同一版本的服务部署在两个K8s集群中:K8s集群1运行正常,但K8s集群2上的所有同类服务持续出现NotLeaderOrFollowerException错误,报错信息为:Received invalid metadata error in produce request on partition TOPIC-4 due to org.apache.kafka.common.errors.NotLeaderOrFollowerException....,消息偶尔能成功发送到1-2个节点,但大部分情况失败。

环境信息

  • SpringBoot 3.1.5(曾升级至3.3.6,问题未解决)
  • JVM 21
  • Kafka 3.4.1(独立部署,不在K8s集群内;曾将客户端升级至3.9.0,无效果)

Kafka集群与Topic分区状态

Kafka节点信息

  • Broker 1:1.1.1.232
  • Broker 2:1.1.1.233
  • Broker 3:1.1.1.234

Topic分区领导者分布(通过Offset Explorer查看)

  • 分区0:1.1.1.232:9092
  • 分区1:1.1.1.233:9092
  • 分区2:1.1.1.234:9092
  • 分区3:1.1.1.234:9092
  • 分区4:1.1.1.232:9092
  • 分区5:1.1.1.233:9092

日志分析

生产者核心日志片段:

[Producer clientId=producer-1] Nodes with data ready to send: [1.1.1.232:9092 (id: 1 rack: null)]

[Producer clientId=producer-1] Sent produce request to 1: (type=ProduceRequest, acks=-1, timeout=30000, partitionRecords=([PartitionProduceData(index=4, records=MemoryRecords(size=1303, buffer=java.nio.HeapByteBuffer[pos=0 lim=1303 cap=1303]))]), transactionalId=''

[Producer clientId=producer-1] Received produce response from node 1 with correlation id 11119

[Producer clientId=producer-1] Got error produce response with correlation id 11119 on topic-partition TOPIC-4, retrying (2147482932 attempts left). Error: NOT_LEADER_OR_FOLLOWER

[Producer clientId=producer-1] Received invalid metadata error in produce request on partition TOPIC-4 due to org.apache.kafka.common.errors.NotLeaderOrFollowerException: For requests intended only for the leader, this error indicates that the broker is not the current leader. For requests intended for any replica, this error indicates that the broker is not a replica of the topic partition.. Going to request metadata update now

[Producer clientId=producer-1] Updating last seen epoch for partition TOPIC-1 from 0 to epoch 0 from new metadata
[Producer clientId=producer-1] Updating last seen epoch for partition TOPIC-2 from 0 to epoch 0 from new metadata
[Producer clientId=producer-1] Updating last seen epoch for partition TOPIC-0 from 0 to epoch 0 from new metadata
[Producer clientId=producer-1] Updating last seen epoch for partition TOPIC-3 from 0 to epoch 0 from new metadata
[Producer clientId=producer-1] Updating last seen epoch for partition TOPIC-5 from 0 to epoch 0 from new metadata
[Producer clientId=producer-1] Updating last seen epoch for partition TOPIC-4 from 0 to epoch 0 from new metadata

[Producer clientId=producer-1] Updated cluster metadata updateVersion 4992 to MetadataSnapshot{clusterId='JBNBUWDLRIOjQHAEqIM2FA', nodes={1=1.1.1.232:9092 (id: 1 rack: null), 2=1.1.1.233:9092 (id: 2 rack: null), 3=1.1.1.234:9092 (id: 3 rack: null)}, partitions=[PartitionMetadata(error=NONE, partition=TOPIC-3, leader=Optional[3], leaderEpoch=Optional[0], replicas=3,1,2, isr=3,1,2, offlineReplicas=), PartitionMetadata(error=NONE, partition=TOPIC-2, leader=Optional[3], leaderEpoch=Optional[0], replicas=3,1,2, isr=3,1,2, offlineReplicas=), PartitionMetadata(error=NONE, partition=TOPIC-5, leader=Optional[2], leaderEpoch=Optional[0], replicas=2,3,1, isr=2,3,1, offlineReplicas=), PartitionMetadata(error=NONE, partition=TOPIC-4, leader=Optional[1], leaderEpoch=Optional[0], replicas=1,2,3, isr=1,2,3, offlineReplicas=), PartitionMetadata(error=NONE, partition=TOPIC-1, leader=Optional[2], leaderEpoch=Optional[0], replicas=2,3,1, isr=2,3,1, offlineReplicas=), PartitionMetadata(error=NONE, partition=TOPIC-0, leader=Optional[1], leaderEpoch=Optional[0], replicas=1,2,3, isr=1,2,3, offlineReplicas=)], controller=1.1.1.234:9092 (id: 3 rack: null)}

[Producer clientId=producer-1] Nodes with data ready to send: [1.1.1.232:9092 (id: 1 rack: null)]
...
[Producer clientId=producer-1] Got error ....

日志关键观察

生产者向Broker 1(1.1.1.232:9092)发送分区4的消息时收到NOT_LEADER_OR_FOLLOWER错误;但请求元数据更新后,元数据明确显示分区4的领导者就是Broker 1,再次发送仍会触发相同错误。偶尔节点重连后,能成功发送的节点会变化,但问题无法彻底解决。

已验证的现象

  • 手动指定已成功发送过消息的分区时,消息能正常投递:
    new ProducerRecord<>("TOPIC", successPartition, uuid, message);
    
  • 服务代码为Kafka官方标准实现,与K8s集群1上的运行版本完全一致,仅按每分钟一次频率发送消息:
    var record = new ProducerRecord<>("TOPIC", uuid, message);  
    SendResult<String, String> result = kafkaTemplate.send(record).get(2000, MILLISECONDS);
    

已执行的排查步骤

  • 执行netstat -anp | grep 9092,确认服务已连接所有Kafka节点;
  • 增加生产者超时时间,问题未解决;
  • 升级Spring Boot至3.3.6、Kafka客户端至3.9.0,无效果;
  • 重建单分区Topic,问题仍存在。

排查方向请求

请从代码、Kubernetes、Kafka、基础设施等维度提供排查建议。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 15:36:06