KafkaJS中consumer.disconnect()无法解析,程序无限阻塞求助
KafkaJS消费到最新偏移量后
disconnect()无限阻塞问题 我基于KafkaJS实现获取指定时间范围内的所有消息,搭建了服务端动作:创建Kafka消费者、聚合现有消息并过滤时间范围外的内容,这是初步务实实现。但消费到最新偏移量后,程序在await consumer.disconnect()处无限阻塞,无任何错误或超时,请问可能的原因是什么?
环境信息
- Node.js: 22.6.0
- Next.js: 14.2.6
- Luxon: ^3.5.0
- KafkaJS: ^2.2.4
实现代码
import { Kafka, EachMessagePayload, KafkaConfig } from 'kafkajs'; import { DateTime } from 'luxon'; export interface MessageEnvelop { topic: string; key: string | null; value: string | null; timestamp: DateTime; } const kafkaConfig = { clientId: `my-client`, brokers: [process.env.BOOTSTRAP_SERVERS!!], ssl: { // @ts-ignore 'ssl.endpoint.identification.algorithm': 'https' }, sasl: { mechanism: 'plain', username: process.env.KAFKA_USERNAME!!, password: process.env.KAFKA_PASSWORD!! } } satisfies KafkaConfig; const client = new Kafka(kafkaConfig); const adminClient = client.admin(); async function getLatestOffsets(topic: string) { await adminClient.connect(); const offsets = await adminClient.fetchTopicOffsets(topic); await adminClient.disconnect(); return offsets; } async function resetOffsets(consumerSuffix: string, topic: string) { console.log(`Resetting offsets for topic ${topic}`); await adminClient.connect(); await adminClient.resetOffsets({ groupId: 'my-consumer-group', topic, earliest: true }) await adminClient.disconnect(); } async function getMessages(consumerSuffix: string, topic: string, fromTime: DateTime, toTime: DateTime, timestampExtractor: (message: string) => string): Promise<MessageEnvelop[]> { await resetOffsets(consumerSuffix, topic); const consumer = client.consumer({ groupId: 'my-consumer-group' }); await consumer.connect(); await consumer.subscribe({ topic, fromBeginning: true }); const latestOffsets = await getLatestOffsets(topic); const messages: MessageEnvelop[] = []; const partitionOffsets: { [partition: number]: number } = {}; latestOffsets.forEach(({ partition, offset }) => { partitionOffsets[partition] = parseInt(offset, 10); }); console.log("Partition offsets: ", partitionOffsets); return new Promise(async (resolve, reject) => { consumer.run({ eachMessage: async ({ topic, partition, message }: EachMessagePayload) => { console.log(`Received message from topic ${topic} on partition ${partition} with offset ${message.offset}`); const messageTimestamp = DateTime.fromISO(timestampExtractor(message.value!!.toString())); if (messageTimestamp >= fromTime && messageTimestamp <= toTime) { messages.push({ topic: topic, key: message.key?.toString() || null, value: message.value?.toString() || null, timestamp: messageTimestamp }); } if (message.offset === (partitionOffsets[partition]-1).toString()) { partitionOffsets[partition] = -1; } console.log(`Partition ${partition} has offset ${message.offset} and latest offset is ${partitionOffsets[partition]}`); console.log(partitionOffsets); if (Object.values(partitionOffsets).every(offset => offset === -1)) { console.log(`consumed ${messages.length} messages. Disconnecting consumer ...`); await consumer.disconnect(); // <---- 阻塞位置 console.log("resolving promise ..."); resolve(messages); } } }) .catch((error) => { console.error("Error while consuming messages: ", error); reject(error); }); }); }
package.json
{ "name": "zip-debugger", "version": "0.1.0", "private": true, "scripts": { "dev": "next dev", "build": "next build", "start": "next start", "lint": "next lint" }, "dependencies": { "@emotion/cache": "^11.13.1", "@emotion/react": "^11.13.3", "@emotion/styled": "^11.13.0", "@mui/material-nextjs": "^5.16.6", "luxon": "^3.5.0", "next": "14.2.6", "react": "^18", "react-dom": "^18", "uuid": "^10.0.0" }, "devDependencies": { "@types/luxon": "^3.4.2", "@types/node": "^20", "@types/react": "^18", "@types/react-dom": "^18", "@types/uuid": "^10.0.0", "eslint": "^8", "eslint-config-next": "14.2.6", "kafkajs": "^2.2.4", "postcss": "^8", "tailwindcss": "^3.4.1", "typescript": "^5" } }
问题原因及解决方案
1. 消费循环未提前终止
KafkaJS的consumer.run()会持续运行消费循环,直到调用consumer.stop()主动终止。直接调用disconnect()时,断开操作会等待消费循环结束,而消费循环默认不会主动停止,导致阻塞。
2. 异步上下文冲突
在eachMessage回调内部调用disconnect(),会和consumer.run()的异步执行上下文产生冲突,断开操作无法正确终止正在运行的消费任务。
修复后的代码
async function getMessages(consumerSuffix: string, topic: string, fromTime: DateTime, toTime: DateTime, timestampExtractor: (message: string) => string): Promise<MessageEnvelop[]> { await resetOffsets(consumerSuffix, topic); const consumer = client.consumer({ groupId: 'my-consumer-group' }); await consumer.connect(); await consumer.subscribe({ topic, fromBeginning: true }); const latestOffsets = await getLatestOffsets(topic); const messages: MessageEnvelop[] = []; const partitionOffsets: { [partition: number]: number } = {}; latestOffsets.forEach(({ partition, offset }) => { partitionOffsets[partition] = parseInt(offset, 10); }); console.log("Partition offsets: ", partitionOffsets); return new Promise((resolve, reject) => { // 封装停止与断开逻辑 const shutdownConsumer = async () => { try { await consumer.stop(); // 先终止消费循环 await consumer.disconnect(); console.log("resolving promise ..."); resolve(messages); } catch (err) { reject(err); } }; consumer.run({ eachMessage: async ({ topic, partition, message }: EachMessagePayload) => { console.log(`Received message from topic ${topic} on partition ${partition} with offset ${message.offset}`); const messageTimestamp = DateTime.fromISO(timestampExtractor(message.value!!.toString())); if (messageTimestamp >= fromTime && messageTimestamp <= toTime) { messages.push({ topic: topic, key: message.key?.toString() || null, value: message.value?.toString() || null, timestamp: messageTimestamp }); } if (message.offset === (partitionOffsets[partition]-1).toString()) { partitionOffsets[partition] = -1; } console.log(`Partition ${partition} has offset ${message.offset} and latest offset is ${partitionOffsets[partition]}`); console.log(partitionOffsets); if (Object.values(partitionOffsets).every(offset => offset === -1)) { console.log(`consumed ${messages.length} messages. Disconnecting consumer ...`); shutdownConsumer(); // 调用封装的停止逻辑 } } }).catch((error) => { console.error("Error while consuming messages: ", error); reject(error); }); }); }
额外优化建议
- 复用AdminClient:当前
getLatestOffsets和resetOffsets每次调用都重新连接AdminClient,可改为复用单例实例,减少连接开销。 - 隔离消费者组:如果多个请求同时调用该函数,使用同一个消费者组会导致偏移量冲突,建议为每个请求生成唯一的消费者组ID,或确保同一时间仅一个实例消费。
- 错误处理增强:在
shutdownConsumer中添加错误捕获,避免断开过程中的异常导致Promise无法resolve/reject。
内容的提问来源于stack exchange,提问作者Peter C. Glade
相关产品推荐
相关产品推荐

