Kafkajs eachMessage无法拉取消息,如何等待全部消息获取?
问题分析
KafkaJS 的 consumer.run() 是一个长期运行的异步函数,它会持续监听 Kafka 主题并触发 eachMessage 回调,不会自动 resolve。你原代码里的逻辑矛盾在于:
- 加
await时,consumer.run()永远不会完成,后续的disconnect和return messages根本执行不到,所以返回空数组; - 去掉
await时,代码会直接跳到disconnect,此时消费者还没来得及拉取消息,只有脚本结束前可能收到少量消息,但无法正确收集。
要实现「拉取所有消息后返回」,必须主动控制消费者的停止时机,以下是几种实用方案:
方案1:拉取指定数量的消息
如果明确需要拉取的消息总数,用计数器触发停止:
const { Kafka, logLevel } = require('kafkajs') async function consume_messages(config, targetCount) { const kafka = new Kafka({ logLevel: logLevel.INFO, brokers: [config.broker], ssl: true, sasl: { mechanism: config.mechanism, // 原代码多了不必要的数组包裹,直接传字符串即可 username: config.Username, password: config.Password }, }) const topic = config.client_id const consumer = kafka.consumer({ groupId: 'my-group', fromBeginning: true }) await consumer.connect(); await consumer.subscribe({ topics: [topic], fromBeginning: true }) let messages = [] const stopConsumer = async () => { await consumer.stop() await consumer.disconnect() } await consumer.run({ eachMessage: async ({ message }) => { messages.push(message) console.log('RECEIVED MESSAGE', JSON.parse(message.value), message.offset); // 达到目标数量后停止消费 if (messages.length >= targetCount) { await stopConsumer() } } }) return messages ; }
方案2:拉取主题所有历史消息(消费到分区末尾)
如果要拉取当前主题的全部历史消息,需要先获取每个分区的最新偏移量,判断是否消费到末尾:
const { Kafka, logLevel } = require('kafkajs') async function consume_messages(config) { const kafka = new Kafka({ logLevel: logLevel.INFO, brokers: [config.broker], ssl: true, sasl: { mechanism: config.mechanism, username: config.Username, password: config.Password }, }) const topic = config.client_id const consumer = kafka.consumer({ groupId: 'my-group', fromBeginning: true }) await consumer.connect(); await consumer.subscribe({ topics: [topic], fromBeginning: true }) // 获取每个分区的最新偏移量 const admin = kafka.admin() await admin.connect() const topicOffsets = await admin.fetchTopicOffsets(topic) await admin.disconnect() // 记录每个分区已消费的最大偏移量 const consumedOffsets = new Map() let messages = [] const stopConsumer = async () => { await consumer.stop() await consumer.disconnect() } await consumer.run({ eachMessage: async ({ message, partition }) => { messages.push(message) console.log('RECEIVED MESSAGE', JSON.parse(message.value), message.offset); // 更新当前分区的已消费偏移量 consumedOffsets.set(partition, parseInt(message.offset)) // 检查所有分区是否都消费到最新偏移量(latestOffset是下一个待写入的偏移量,所以要减1) const allCompleted = Array.from(topicOffsets).every(([partition, latestOffset]) => { return consumedOffsets.get(partition) === latestOffset - 1 }) if (allCompleted) { await stopConsumer() } } }) return messages ; }
方案3:超时自动停止
如果不确定消息数量,可设置超时时间,超时后停止消费:
const { Kafka, logLevel } = require('kafkajs') async function consume_messages(config, timeoutMs = 5000) { const kafka = new Kafka({ logLevel: logLevel.INFO, brokers: [config.broker], ssl: true, sasl: { mechanism: config.mechanism, username: config.Username, password: config.Password }, }) const topic = config.client_id const consumer = kafka.consumer({ groupId: 'my-group', fromBeginning: true }) await consumer.connect(); await consumer.subscribe({ topics: [topic], fromBeginning: true }) let messages = [] const stopConsumer = async () => { await consumer.stop() await consumer.disconnect() } // 超时后自动停止消费 setTimeout(async () => { await stopConsumer() }, timeoutMs) await consumer.run({ eachMessage: async ({ message }) => { messages.push(message) console.log('RECEIVED MESSAGE', JSON.parse(message.value), message.offset); } }) return messages ; }
额外注意
原代码有两个语法/逻辑错误:
sasl.mechanism不需要用数组包裹,直接传对应机制的字符串(如'plain'、'scram-sha-256')即可;consumer.run()末尾多了一个多余的}),需删除。
内容的提问来源于stack exchange,提问作者Gil Allen
相关产品推荐
相关产品推荐

