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

Node.js Kafka普通消费者正常收消息,消费者组请求超时

问题描述

本地Kafka集群已创建test_producer主题,通过以下命令启动的控制台生产者可正常工作:

bin/kafka-console-producer.sh --topic test_producer --bootstrap-server localhost:9092

Node.js应用使用kafka-node库时,普通消费者能正常接收消息,但创建消费者组时抛出请求超时异常。消费者组代码如下:

import { ConsumerGroup } from "kafka-node";
const options = {
    kafkaHost: "127.0.0.1:9092",
    batch: undefined,
    ssl: true,
    groupId: "ExampleGroup",
    sessionTimeout: 15000,
    protocol: ["roundrobin"],
    encoding: "utf8",
    fromOffset: "latest", // default
    commitOffsetsOnFirstJoin: true,
    outOfRangeOffset: "earliest",
    onRebalance: (isAlreadyMember, callback) => {
      callback();
    },
  };
  const consumerGroup = new ConsumerGroup(options, ["test_producer"]);
  consumerGroup.on("message", (message) => {
    console.log(`group:${message}`);
  });

Kafka的config/server.properties配置如下:

############################# Server Basics #############################
broker.id=0

############################# Socket Server Settings #############################

num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600

############################# Log Basics #############################

log.dirs=/home/kafka/logs
num.partitions=1
num.recovery.threads.per.data.dir=1

############################# Internal Topic Settings  #############################
offsets.topic.replication.factor=1
transaction.state.log.replication.factor=1
transaction.state.log.min.isr=1

############################# Log Flush Policy #############################

#log.flush.interval.messages=10000
#log.flush.interval.ms=1000

############################# Log Retention Policy #############################

log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

############################# Zookeeper #############################
zookeeper.connect=localhost:2181
zookeeper.connection.timeout.ms=18000


############################# Group Coordinator Settings #############################

group.initial.rebalance.delay.ms=0
delete.topic.enable = true
advertised.listeners=PLAINTEXT://127.0.0.1:9092
listeners = PLAINTEXT://127.0.0.1:9092

问题分析

  1. SSL配置不匹配:消费者组开启了ssl: true,但Kafka集群仅配置了PLAINTEXT无加密监听,导致连接握手失败超时。普通消费者未启用SSL,因此能正常工作。
  2. 库兼容性问题:kafka-node库维护状态停滞,与较新版本Kafka的消费者组协议可能存在兼容性问题。
  3. 超时参数设置过短:sessionTimeout=15000(15秒)可能因集群初始化或网络延迟导致超时。

解决方案

1. 修正SSL配置

将消费者组配置中的ssl改为false,匹配Kafka的PLAINTEXT协议:

const options = {
    // ...其他配置
    ssl: false, // 关闭SSL
    // ...其他配置
};

2. 调整超时参数

增大sessionTimeout至30秒,降低超时概率:

sessionTimeout: 30000,

3. 检查内部主题状态

消费者组依赖__consumer_offsets主题存储偏移量,检查该主题状态:

bin/kafka-topics.sh --describe --topic __consumer_offsets --bootstrap-server localhost:9092

若主题异常,可删除后重启Kafka(会丢失所有消费者组偏移量):

bin/kafka-topics.sh --delete --topic __consumer_offsets --bootstrap-server localhost:9092

4. 替换为活跃的Kafka客户端

kafka-node已停止维护,建议改用kafkajs,消费者组示例代码如下:

import { Kafka } from 'kafkajs'

const kafka = new Kafka({
  clientId: 'example-group-client',
  brokers: ['localhost:9092']
})

const consumer = kafka.consumer({ groupId: 'ExampleGroup' })

const run = async () => {
  await consumer.connect()
  await consumer.subscribe({ topic: 'test_producer', fromBeginning: false })

  await consumer.run({
    eachMessage: async ({ topic, partition, message }) => {
      console.log(`group: ${message.value.toString()}`)
    },
  })
}

run().catch(console.error)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 11:55:26