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
问题分析
- SSL配置不匹配:消费者组开启了
ssl: true,但Kafka集群仅配置了PLAINTEXT无加密监听,导致连接握手失败超时。普通消费者未启用SSL,因此能正常工作。 - 库兼容性问题:
kafka-node库维护状态停滞,与较新版本Kafka的消费者组协议可能存在兼容性问题。 - 超时参数设置过短:
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
相关产品推荐
相关产品推荐

