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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 04:40:29