如何在ExpressJS有状态API端点中复用单个Kafka Consumer?
解决方案
1. 重构KafkaConfig:单例消费者+消息广播机制
首先修改KafkaConfig.js,让消费者只初始化并启动一次,同时提供注册/注销消息监听器的能力,实现消息向所有关注客户端的分发:
const { Kafka, Partitioners } = require('kafkajs'); class KafkaConfig { constructor() { this.kafka = new Kafka({ clientId: "BOOKING_SERVICE", brokers: ['127.0.0.1:39092', '127.0.0.1:29092'], }); this.producer = this.kafka.producer({ createPartitioner: Partitioners.LegacyPartitioner }); this.consumer = this.kafka.consumer({ groupId: "test-group" }); this.messageListeners = new Set(); // 存储所有消息监听器 // 初始化时自动启动消费者 this.initConsumer(); } async produce(messages, topic) { try { await this.producer.connect(); await this.producer.send({ topic: topic, messages: [{ value: messages }], }); } catch (error) { console.error(error); } finally { await this.producer.disconnect(); } } async initConsumer() { try { await this.consumer.connect(); await this.consumer.subscribe({ topics: ['driver_history', 'Booking_Modal', 'Drivers_stats', 'DRIVER_GPS_DATA'], fromBeginning: true }); await this.consumer.run({ eachMessage: async ({ topic, partition, message }) => { const value = JSON.parse(message.value); // 把消息广播给所有注册的监听器 this.messageListeners.forEach(listener => listener(value)); }, }); console.log("Kafka consumer initialized and running"); } catch (error) { console.error("Kafka consumer initialization failed:", error); } } // 注册消息监听器 addMessageListener(listener) { this.messageListeners.add(listener); } // 注销消息监听器 removeMessageListener(listener) { this.messageListeners.delete(listener); } } // 导出单例实例,确保整个应用只有一个KafkaConfig实例 module.exports = new KafkaConfig();
2. 实现SSE客户端管理逻辑
在路由文件中维护活跃客户端集合,关联每个客户端的bookingId与响应对象,同时处理客户端断开后的清理工作:
const express = require('express'); const router = express.Router(); const kafkaConfig = require('./KafkaConfig'); // 导入单例的KafkaConfig // 存储活跃SSE客户端:key为bookingId,value为对应响应对象的集合 const activeClients = new Map(); // 统一处理Kafka消息的函数 const handleKafkaMessage = (message) => { if (!message.documentKey || !message.documentKey._id) return; const bookingId = message.documentKey._id; // 针对匹配的bookingId推送消息 if (activeClients.has(bookingId)) { const clients = activeClients.get(bookingId); clients.forEach(res => { try { res.write(`data: ${JSON.stringify(message)}\n\n`); } catch (err) { // 客户端连接异常,移除该客户端 clients.delete(res); if (clients.size === 0) { activeClients.delete(bookingId); } } }); } }; // 仅注册一次Kafka消息处理函数 kafkaConfig.addMessageListener(handleKafkaMessage); router.get("/bookingSSE", (req, res) => { const bookingId = req.query.bookingId; // 假设bookingId从查询参数获取,可根据实际场景调整 if (!bookingId) { res.status(400).send("Error: bookingId parameter is missing."); return; } // 设置SSE响应头 res.setHeader("Content-Type", "text/event-stream"); res.setHeader("Cache-Control", "no-cache"); res.setHeader("Connection", "keep-alive"); res.setHeader("Access-Control-Allow-Origin", "*"); // 根据实际需求调整跨域配置 // 发送连接成功提示 res.write(`data: Connected!!\n\n`); // 将当前客户端加入活跃列表 if (!activeClients.has(bookingId)) { activeClients.set(bookingId, new Set()); } const clientSet = activeClients.get(bookingId); clientSet.add(res); // 处理客户端主动断开 req.on("close", () => { console.log(`Client for bookingId ${bookingId} disconnected.`); clientSet.delete(res); if (clientSet.size === 0) { activeClients.delete(bookingId); } }); // 处理连接异常 req.on("error", (err) => { console.error(`Error for bookingId ${bookingId}:`, err); clientSet.delete(res); if (clientSet.size === 0) { activeClients.delete(bookingId); } }); }); module.exports = router;
关键改进说明
- 单例消费者:通过导出KafkaConfig单例,确保应用生命周期内只有一个Kafka消费者连接,避免重复连接引发的错误。
- 消息广播:消费者接收到消息后,通过统一的监听器分发给所有匹配的客户端,无需为每个客户端单独启动消费逻辑。
- 客户端生命周期管理:维护活跃客户端集合,在客户端断开时及时清理,避免内存泄漏和无效消息推送。
- 避免重复注册:仅在路由初始化时注册一次消息处理函数,不会因多次请求重复添加监听器。
内容的提问来源于stack exchange,提问作者0xSurya
相关产品推荐
相关产品推荐

