使用KafkaJS连接Kafka消费者时遇主题分区非领导者等错误
KafkaJS消费者报错:"This server is not the leader for that topic-partition"及"Specified group generation id is not valid"
消费者配置代码
private createKafkaConsumer(): Consumer { const kafka = new Kafka({ logLevel: logLevel.WARN, brokers: this.kafkaHost, connectionTimeout: 10000, requestTimeout: 25000, retry: { initialRetryTime: 100, retries: 8, }, }); const consumer = kafka.consumer({ groupId: this.groupId, sessionTimeout: 60000, rebalanceTimeout: 60000, heartbeatInterval: 3000, allowAutoTopicCreation: false, maxWaitTimeInMs: 5000, maxBytes: 1024 * 10 }); return consumer; }
错误日志
{"level":"ERROR","timestamp":"2022-10-16T12:29:57.236Z","logger":"kafkajs","message":"[Connection] Response Fetch(key: 1, version: 10)","broker":"b-1.ccss-prod-pulse-msk.5jyfxf.c3.kafka.ap-south-1.amazonaws.com:9092","clientId":"kafkajs","error":"This server is not the leader for that topic-partition","correlationId":2301,"size":162} {"level":"ERROR","timestamp":"2022-10-16T12:29:57.258Z","logger":"kafkajs","message":"[Connection] Response Fetch(key: 1, version: 10)","broker":"b-1.ccss-prod-pulse-msk.5jyfxf.c3.kafka.ap-south-1.amazonaws.com:9092","clientId":"kafkajs","error":"This server is not the leader for that topic-partition","correlationId":2302,"size":162} {"level":"ERROR","timestamp":"2022-10-16T12:29:57.275Z","logger":"kafkajs","message":"[Connection] Response Fetch(key: 1, version: 10)","broker":"b-1.ccss-prod-pulse-msk.5jyfxf.c3.kafka.ap-south-1.amazonaws.com:9092","clientId":"kafkajs","error":"This server is not the leader for that topic-partition","correlationId":2303,"size":162} {"level":"ERROR","timestamp":"2022-10-16T12:29:57.331Z","logger":"kafkajs","message":"[Connection] Response OffsetCommit(key: 8, version: 5)","broker":"b-2.ccss-prod-pulse-msk.5jyfxf.c3.kafka.ap-south-1.amazonaws.com:9092","clientId":"kafkajs","error":"Specified group generation id is not valid","correlationId":6017,"size":86} {"level":"ERROR","timestamp":"2022-10-16T12:29:57.334Z","logger":"kafkajs","message":"[Connection] Response Fetch(key: 1, version: 10)","broker":"b-1.ccss-prod-pulse-msk.5jyfxf.c3.kafka.ap-south-1.amazonaws.com:9092","clientId":"kafkajs","error":"This server is not the leader for that topic-partition","correlationId":2304,"size":162} {"level":"ERROR","timestamp":"2022-10-16T12:29:59.872Z","logger":"kafkajs","message":"[Consumer] Crash: KafkaJSNonRetriableError: Specified group generation id is not valid","groupId":"in-EW-Maintenance-Hourmeter.1","stack":"KafkaJSNonRetriableError: Specified group generation id is not valid\n at /usr/src/app/node_modules/kafkajs/src/retry/index.js:55:18\n at runMicrotasks (<anonymous>)\n at processTicksAndRejections (node:internal/process/task_queues:94:5)"}
解决方向
处理分区leader变更
这个报错说明消费者连接的broker不是目标分区的leader,大概率是集群发生了leader选举或分区重分配。先检查Kafka集群状态,确认所有分区leader是否正常;然后重启消费者实例,让KafkaJS重新拉取最新的分区元数据。修复消费组generation id无效问题
该错误通常是因为消费组的generation已过期,比如消费者心跳超时被踢出组,或者同组内其他消费者触发了重平衡。- 检查
sessionTimeout和heartbeatInterval:当前配置符合Kafka推荐(心跳间隔小于会话超时的1/3),但如果集群网络不稳定,可尝试将sessionTimeout调至90000ms; - 确认消费组内无残留的旧实例,若有异常退出的实例,可手动清理消费组偏移记录(测试环境或确认无数据丢失风险时操作):
kafka-consumer-groups.sh --bootstrap-server <broker地址> --delete --group in-EW-Maintenance-Hourmeter.1 - 确保消费者关闭时调用
consumer.disconnect(),避免异常退出导致状态残留。
- 检查
优化KafkaJS配置
- 添加
metadataMaxAge参数,缩短元数据刷新间隔,加快集群变化感知:const kafka = new Kafka({ // 其他配置 metadataMaxAge: 30000 // 30秒刷新一次元数据 }); - 增加
retry配置的maxRetryTime,应对短暂网络波动:retry: { initialRetryTime: 100, retries: 8, maxRetryTime: 30000 }
- 添加
内容的提问来源于stack exchange,提问作者Mohamed Anser Ali
相关产品推荐
相关产品推荐

