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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 00:46:09