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

如何在Express.js服务器中获取所有已建立的HTTP连接(SSE场景)

解决Express.js SSE连接的统一管理问题

核心问题分析

你当前的代码尝试用this.connections.push(req)存储连接,但存在几个关键问题:

  • 存储req而非包含req/res的完整连接对象,后续无法直接通过连接推送消息
  • 连接关闭时未从数组中移除对应项,会导致数组堆积无效连接
  • streamCount的维护时机错误,没有在新连接建立时完成递增操作

修正后的完整代码

const { EventEmitter } = require('events');
const express = require('express');
const sseRouter = express.Router();

// 全局SSE实例
const sse = new SSE();

sseRouter.get("/stream", (req, res) => {
   sse.init(req, res);
});

class SSE extends EventEmitter {    
  constructor() {
    super();
    this.connections = []; // 存储有效连接的数组
    this.streamCount = 0;  // 将计数移到类内部,统一管理
  }

  init(req, res) {
    res.writeHead(200, {
      Connection: "keep-alive",
      "Content-Type": "text/event-stream",
      "Cache-Control": "no-cache",
    });
    console.log("client connected init...");

    // 创建连接对象,同时存储req和res,方便后续操作
    const connection = { req, res };
    this.connections.push(connection);
    this.streamCount++; // 新连接建立,计数递增

    let id = 0;
    // 发送初始消息
    res.write(`data: some data \n\n`);
   
    const dataListener = (data) => {
      // 修正原代码格式错误:避免重复输出event字段,data字段才是消息内容载体
      if (data.event) {
        res.write(`event: ${data.event} \n`);
      }
      res.write(`data: ${data.data} \n`);
      res.write(`id: ${++id} \n`);
      res.write("\n");
    };

    this.on("data", dataListener);

    req.on("close", () => {
      this.removeListener("data", dataListener);
      // 从数组中移除当前无效连接
      this.connections = this.connections.filter(conn => conn !== connection);
      this.streamCount--; // 连接关闭,计数递减
      console.log("Stream closed, current connections:", this.streamCount);
    });        
  }

  // 可选:批量推送方法,给所有在线客户端发消息
  broadcast(data) {
    this.connections.forEach(conn => {
      try {
        conn.res.write(`data: ${JSON.stringify(data)} \n\n`);
      } catch (err) {
        // 处理连接已断开的异常
        console.error("Failed to send to connection:", err);
      }
    });
  }
}; 

module.exports = sseRouter;

关键改进点

  • 统一存储连接对象:将包含req和res的对象存入this.connections,既能标识唯一连接,也能直接通过res推送消息
  • 正确维护连接列表:连接关闭时用filter从数组中移除无效连接,避免内存泄漏
  • 计数与列表同步:把streamCount移到类内部,和连接列表的增删操作绑定,确保计数准确
  • 修正SSE消息格式:修复原代码重复输出event字段的错误,符合SSE规范
  • 新增广播方法:可选的broadcast方法可以快速给所有在线客户端推送消息

验证唯一连接的方式

每个EventSource实例会建立独立的HTTP长连接,你可以通过以下方式验证:

  • 在init方法中打印req.socket.remoteAddress和req.socket.remotePort,每个连接的端口不同,代表唯一连接
  • 打开多个浏览器标签页,观察this.streamCount和this.connections.length是否同步递增

内容的提问来源于stack exchange,提问作者Александр Сосо

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 19:40:28