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

3节点Kafka集群停Leader后Node.js生产者持续报分区无Leader错误

问题原因与解决方案

核心原因分析

  1. 目标Topic副本数不足:若test-topic创建时未指定副本数,默认仅为1。当Leader Broker下线后,无其他同步副本可成为新Leader,导致领导选举无法完成,持续报错。
  2. 生产者未配置重试逻辑:领导选举需要短暂时间完成,当前代码未对发送操作配置重试,一旦遭遇选举中错误就直接抛出,未等待选举完成后重试。
  3. 异步错误未捕获:setInterval内的异步发送错误未被捕获,可能导致后续发送逻辑异常中断。

解决步骤

1. 确保Topic拥有3个副本

先检查test-topic的副本配置:

# 进入任意Kafka容器执行
docker exec -it kafka1 kafka-topics --describe --topic test-topic --bootstrap-server kafka1:9092

若ReplicationFactor为1,需重新创建Topic(或修改现有Topic):

# 删除原有Topic(若允许)
docker exec -it kafka1 kafka-topics --delete --topic test-topic --bootstrap-server kafka1:9092

# 创建带3个副本的Topic
docker exec -it kafka1 kafka-topics --create --topic test-topic --partitions 1 --replication-factor 3 --bootstrap-server kafka1:9092

2. 配置生产者重试机制

修改Node.js生产者代码,添加重试配置并捕获发送错误:

console.log("producer..........")
const { Kafka } = require('kafkajs')

const kafka = new Kafka({
  clientId: 'my-app',
  brokers: ['localhost:8092', 'localhost:8093', 'localhost:8094']
})

// 开启生产者重试,适配领导选举窗口
const producer = kafka.producer({
  retry: {
    retries: 10,
    initialRetryTime: 1000,
    factor: 2,
    maxRetryTime: 30000
  }
})

const run = async () => {
  await producer.connect()
  setInterval(async () => {
    try {
      await producer.send({
        topic: 'test-topic',
        messages: [{ value: 'Hello KafkaJS user!' }],
      })
    } catch (err) {
      console.error("发送失败,等待重试:", err.message)
    }
  }, 2000)
}

run().catch(console.error)

3. 验证集群副本同步状态

检查Broker日志确认副本同步正常:

docker logs kafka2 | grep "ISR"

正常情况下,test-topic的ISR(同步副本集)应包含所有3个Broker节点。


额外检查点

  • 确认ZooKeeper容器正常运行:使用docker ps查看状态。
  • 关闭Leader Broker后,等待10-20秒让选举完成,再观察生产者是否恢复发送。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 01:17:45