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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 08:52:34