NodeJS中Kafka消费者跨Topic消费异常问题求助
问题根源分析
你的问题出在KafkaConsumer单例复用以及连接状态的错误判断上:
- 两次调用
startReadingFromA1和startReadingFromA2时,都是通过KafkaConsumer.getInstance()获取同一个单例实例 - 第一次初始化A1消费者时,
connectToKafkaCluster会将this.isConnected设为true - 初始化A2消费者时,因为
this.isConnected已经是true,跳过了创建新消费者连接的步骤,直接调用subscribeConsumer把A2的Topic也订阅到了同一个消费者实例上 - 两次调用
this.getMessage相当于给同一个消费者注册了两个回调函数,导致这个消费者同时消费A1和A2的消息,并且两个回调都会处理所有消息,最终表现为A1的消费者逻辑在处理A2的消息
修复方案
需要让每个Topic对应独立的Kafka消费者实例,不能复用同一个单例的状态。具体修改点:
- 取消KafkaConsumer的单例模式,每次初始化消费者都创建新的实例
- 移除
this.isConnected全局状态判断,每个消费者实例自己管理连接状态 - 确保每个消费者实例只订阅对应的单个Topic,并且绑定对应的回调函数
修改后的代码示例
重构KafkaConsumer类(去掉单例)
class KafkaConsumer { constructor() { this.clientId = ''; this.brokerList = []; this.sasl = false; this.ssl = false; this.logLevel = null; this.isConnected = false; this.kafkaConsumerObj = null; } setClientId(clientId) { this.clientId = clientId; } setBrokers(brokerList) { this.brokerList = brokerList; } setSasl(sasl) { this.sasl = sasl; } setSSL() { this.ssl = true; // 根据实际业务逻辑调整 } setLogLevel(logLevel) { this.logLevel = logLevel; } getBrokerList() { return process.env.KAFKA_BROKERS.split(','); // 替换为你的broker列表获取逻辑 } initialize() { return new Kafka({ clientId: this.clientId, brokers: this.brokerList, sasl: this.sasl, ssl: this.ssl, logLevel: logLevel.ERROR, connectionTimeout: DEFAULT_CONNECTION_TIMEOUT, retry: { initialRetryTime: 100, retries: 3, }, }); } async connectToKafkaCluster(kafkaConsumerObj, groupId) { this.kafkaConsumerObj = kafkaConsumerObj.consumer({ groupId }); await this.kafkaConsumerObj.connect(); this.isConnected = true; } async subscribeConsumer(topicName) { await this.kafkaConsumerObj.subscribe({ topic: topicName }); } getMessage(logger, callback, country) { this.kafkaConsumerObj.run({ eachMessage: async ({ message }) => { callback(message, country); }, }); } async initializeConsumer( logger, clientId, topicName, callback, maxBytes = DEFAULT_KAFKA_CONSUMER_MAX_BYTES ) { const country = topicName.split('-')[0]; logger.info(`Initializing Consumer for clientId: ${clientId} and topicName: ${topicName}`); this.setClientId(clientId); let brokerList = this.getBrokerList(); this.setBrokers(brokerList); this.setSasl(false); this.setSSL(); this.setLogLevel(logLevel.ERROR); const kafkaConsumerObj = this.initialize(); const groupId = `group-${topicName.toLowerCase()}-${process.env.NODE_ENV}`; await this.connectToKafkaCluster(kafkaConsumerObj, groupId); await this.subscribeConsumer(topicName); logger.info(`Connected Consumer Kafka cluster for topic ${topicName}`); this.getMessage(logger, callback, country); } }
修改startReading方法(每次创建新实例)
const startReadingFromA1 = async () => { const logger = LogHelper.getLoggerWithoutReq(); try { logger.info(`Initializing the Kafka read from Topic A1`); const Consumer = new KafkaConsumer(); const consumerId = nanoid(); let countryInCaps = process.env.COUNTRY ? process.env.COUNTRY.toUpperCase() : null; await Consumer.initializeConsumer(logger, `CONS-${consumerId}`, `${countryInCaps}-ModuleA`, handleMessages); return true; } catch (err) { logger.error(`Error in startReadingFromA1: ${JSON.stringify(err)}`); } } const startReadingFromA2 = async () => { const logger = LogHelper.getLoggerWithoutReq(); try { logger.info(`Initializing the Kafka read from Topic A2`); const Consumer = new KafkaConsumer(); const consumerId = nanoid(); let countryInCaps = process.env.COUNTRY ? process.env.COUNTRY.toUpperCase() : null; await Consumer.initializeConsumer(logger, `CONS-${consumerId}`, `${countryInCaps}-ModuleB`, handleFailedMessages); return true; } catch (err) { logger.error(`Error in startReadingFromA2: ${JSON.stringify(err)}`); } }
额外注意点
- 原代码中
startReadingFromA1和startReadingFromA2的catch块存在变量名错误(用了未定义的error),修改后的代码已修正 - 确保每个消费者的
groupId唯一,你的代码中通过topicName生成的逻辑是合理的 - 检查
getMessage的消息处理逻辑,确保回调只处理当前消费者订阅的Topic消息
内容的提问来源于stack exchange,提问作者Prabhat Mishra
相关产品推荐
相关产品推荐

