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

如何使用KafkaJS在JavaScript中将Kafka消息存入数组?

问题原因

KafkaJS的consumer.run()方法返回的Promise不会等待所有消息消费完成才resolve,它仅在消费者成功启动消费循环后就立即完成。你在await consumer.run()之后立刻打印arr,此时eachMessage的回调函数还未处理任何消息,所以数组为空。

解决方案

根据你的业务需求,分两种场景处理:

场景1:消费指定数量消息后停止,再使用数组

如果需要消费一批消息(比如所有历史消息)后停止,再处理数组,可以通过计数器控制消费流程:

const consume = async () => {
  await consumer.connect();
  await consumer.subscribe({ topic, fromBeginning: true });

  const arr = [];
  const targetCount = 10; // 可根据需求调整,设为Infinity则持续消费直到手动停止

  await new Promise((resolve) => {
    consumer.run({
      eachMessage: async ({ message }) => {
        // 注意:message.value是Buffer类型,需先转字符串再处理
        const clientData = JSON.stringify(message.value.toString());
        arr.push(clientData);
        
        // 达到目标数量后停止消费
        if (arr.length >= targetCount) {
          await consumer.disconnect();
          resolve();
        }
      },
    });
  });

  console.log(arr); // 此时数组已填充消息
};

场景2:持续消费,实时维护数组供应用全局使用

如果需要持续消费并让应用其他模块随时访问消息数组,可将数组放在模块级作用域:

// consume.js
const { Kafka, logLevel } = require("kafkajs");

const clientId = "my-client";
const brokers = ["localhost:443"];
const topic = "client.attributes";

const kafka = new Kafka({
  clientId,
  brokers,
  logLevel: logLevel.INFO,
  connectionTimeout: 5000,
  retry: {
    initialRetryTime: 3000,
    retries: 2,
  },
});

const consumer = kafka.consumer({
  groupId: clientId,
  minBytes: 5,
  maxBytes: 1e6,
  maxWaitTimeInMs: 3000,
});

// 模块级数组,供外部模块访问
const messageArray = [];

const consume = async () => {
  await consumer.connect();
  await consumer.subscribe({ topic, fromBeginning: true });

  await consumer.run({
    eachMessage: ({ message }) => {
      const clientData = JSON.stringify(message.value.toString());
      messageArray.push(clientData);
      console.log("新增消息,当前数组长度:", messageArray.length);
    },
  });

  console.log("消费者已启动,开始接收消息");
};

module.exports = { consume, messageArray };

在index.js中调用并使用数组:

const { consume, messageArray } = require("./consume");

consume().catch((err) => {
  console.error("error in consumer", err);
});

// 示例:定时打印数组内容,模拟应用其他部分使用消息
setInterval(() => {
  console.log("当前数组内容:", messageArray);
}, 5000);
额外注意事项
  • message.value是Buffer类型,直接JSON.stringify会得到Buffer的序列化结构(如{"type":"Buffer","data":[...]}),需先调用.toString()转换为字符串再处理;如果消息本身是JSON格式,可直接解析:const clientData = JSON.parse(message.value.toString());。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:05:42