3节点Kafka集群停Leader后Node.js生产者持续报分区无Leader错误
问题原因与解决方案
核心原因分析
- 目标Topic副本数不足:若
test-topic创建时未指定副本数,默认仅为1。当Leader Broker下线后,无其他同步副本可成为新Leader,导致领导选举无法完成,持续报错。 - 生产者未配置重试逻辑:领导选举需要短暂时间完成,当前代码未对发送操作配置重试,一旦遭遇选举中错误就直接抛出,未等待选举完成后重试。
- 异步错误未捕获:
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
相关产品推荐
相关产品推荐

