KafkaJS生产者与消费者运行正常但无输出求助
Kafkajs生产者消费者无输出问题排查与解决
问题描述
使用Node.js的Kafkajs模块编写了生产者和消费者代码,但运行后既看不到生产者生成的随机数输出,也看不到消费者获取的数据,需要排查解决。
生产者代码
const {Kafka} = require("kafkajs") const kafka = new Kafka({ clientId : 'my-producer', brokers : ['kafka1:9092'] }) const producer = kafka.producer(); const produceMessage = async () => { const value = Math.random() * 1; console.log(value); try{ await producer.send({ topic : 'tester', messages : [ { value }, ], }) }catch(error) { console.log(error); } } const run = async () => { await producer.connect() setInterval(produceMessage,1000) } run()
消费者代码
const {Kafka} = require("kafkajs") const kafka = new Kafka({ clientId : 'my-consumer', brokers : ['kafka1:9092'] }) const consumer = kafka.consumer({groupId : 'consumer-group'}); const run = async () => { await consumer.connect() await consumer.subscribe({ topic: 'tester' }) await consumer.run({ eachMessage: async ({ topic, message }) => { console.log({ offset: message.offset, value: message.value, }) }, }) } run()
排查与解决步骤
1. 验证Kafka集群可达性
- 确认
brokers配置中的地址正确:如果是本地Kafka实例,通常应为localhost:9092;Docker环境下需确认容器名称/IP是否能被Node进程访问 - 在生产者和消费者的连接逻辑中添加错误捕获,打印完整错误信息:
// 生产者run函数修改 const run = async () => { try { await producer.connect() console.log('生产者连接成功') setInterval(produceMessage, 1000) } catch (error) { console.error('生产者连接失败:', error) } } // 消费者run函数修改 const run = async () => { try { await consumer.connect() console.log('消费者连接成功') await consumer.subscribe({ topic: 'tester' }) await consumer.run({ /* ... */ }) } catch (error) { console.error('消费者连接失败:', error) } }
2. 确保主题tester存在
Kafka默认不会自动创建主题(除非集群开启auto.create.topics.enable配置),可通过两种方式解决:
- 用Kafka命令行工具手动创建主题:
kafka-topics.sh --create --topic tester --bootstrap-server kafka1:9092 --partitions 1 --replication-factor 1 - 在Kafkajs生产者配置中开启自动创建主题:
const producer = kafka.producer({ createTopics: true });
3. 修复消息序列化问题
Kafkajs要求消息值为Buffer、字符串或null,直接传递数字会导致序列化异常,修改生产者的消息生成逻辑:
const produceMessage = async () => { const value = (Math.random() * 1).toString(); // 转为字符串 console.log(value); try { await producer.send({ topic: 'tester', messages: [{ value }], }) } catch(error) { console.error('生产消息失败:', error) } }
4. 调整消费者消费策略
- 订阅主题时添加
fromBeginning: true,确保消费者能获取启动前发送的消息:await consumer.subscribe({ topic: 'tester', fromBeginning: true }) - 消费者输出时将
message.value转为字符串,避免打印Buffer对象:eachMessage: async ({ topic, message }) => { console.log({ offset: message.offset, value: message.value.toString(), }) }
5. 检查消费者分组配置
确保groupId: 'consumer-group'未被其他消费者实例占用,若同一分组有多个消费者,消息会被分摊,可能导致当前实例接收不到消息。
内容的提问来源于stack exchange,提问作者Suraj Joshi
相关产品推荐
相关产品推荐

