如何使用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
相关产品推荐
相关产品推荐

