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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 01:15:38