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

Express中EventEmitter触发事件无响应,如何实现SSE订阅通知?

问题描述

我正在开发一款Web应用,需求是当用户发布帖子时,所有订阅者都能收到通知。我使用Node.js的SSE(Server-Sent Events)来实现,示例代码如下:

postEvent.js

const EventEmitter = require("node:events");
class PostEvent extends EventEmitter {}
module.exports = new PostEvent();

post.js

const PostEvent = require('./postEvent');
app.post('/post', (req, res) => {
   // 从请求中获取数据并存入数据库
  PostEvent.emit('posted');
})

sse.js

const PostEvent = require('./postEvent');
app.get('/sse', (req, res) => {
   res.set('Content-Type', 'text/event-stream'); // 用于无限流式传输服务器响应
   setInterval(() => res.write(`data: my sse data every second\n\n`), 1000); // 此代码正常工作
   PostEvent.on('posted', () => res.write('data: posted\n\n')); // 此代码无响应
})

问题在于,当我向/post路由发送fetch请求时,PostEvent.on('posted', callback)并未执行。请问是我的操作有误,还是该事件系统无法这样工作?有没有其他方法可以实现该功能?

注意:未提及Express模板代码,客户端使用EventSource接收流的功能正常。

补充:这种事件系统能否在HTTP或SSE环境中使用?我并未向/sse发送请求,却想执行其中的函数。


问题分析与解决方案

问题原因

  1. 回调未被注册:只有当客户端发起/sse请求时,sse.js里的PostEvent.on('posted', ...)才会执行绑定。如果没有任何客户端连接到/sse,这个回调根本不存在,事件触发时自然没有响应。
  2. 重复绑定隐患:每次新的/sse连接进来,都会重复绑定posted事件,导致后续一个事件触发时,多个相同回调执行,还可能引发内存泄漏。

修正方案:维护活跃SSE连接列表

正确的做法是把所有活跃的SSE连接统一管理,当posted事件触发时,遍历所有连接推送消息,而不是在每个连接里单独绑定事件。

修正后的代码:

  1. postEvent.js(不变)
const EventEmitter = require("node:events");
class PostEvent extends EventEmitter {}
module.exports = new PostEvent();
  1. sse.js:维护连接数组,处理连接关闭
const PostEvent = require('./postEvent');
// 存储所有活跃的SSE连接
const activeConnections = [];

app.get('/sse', (req, res) => {
  res.set({
    'Content-Type': 'text/event-stream',
    'Cache-Control': 'no-cache',
    'Connection': 'keep-alive'
  });

  // 将当前连接加入列表
  activeConnections.push(res);

  // 处理连接关闭,从列表移除,避免内存泄漏
  req.on('close', () => {
    const index = activeConnections.indexOf(res);
    if (index !== -1) {
      activeConnections.splice(index, 1);
    }
  });

  // 定时发送心跳,避免连接被网关断开
  const heartbeatInterval = setInterval(() => {
    res.write(`data: heartbeat\n\n`);
    // 确保消息立即发送(部分Express版本需要)
    res.flush?.();
  }, 30000);

  // 连接关闭时清除定时任务
  req.on('close', () => {
    clearInterval(heartbeatInterval);
  });
});

// 全局监听posted事件,向所有活跃连接推送消息
PostEvent.on('posted', (postData) => {
  // 遍历所有连接发送通知
  activeConnections.forEach(res => {
    // 携带帖子数据,客户端解析后展示
    res.write(`data: ${JSON.stringify({type: 'post', data: postData})}\n\n`);
    res.flush?.();
  });
});
  1. post.js:触发事件时携带帖子数据
const PostEvent = require('./postEvent');
app.post('/post', (req, res) => {
  // 从请求中获取数据并存入数据库
  const postData = req.body;
  // 存入数据库逻辑...

  // 触发事件时携带帖子数据
  PostEvent.emit('posted', postData);
  res.status(200).json({success: true});
})

补充问题解答

这种事件系统可以在HTTP/SSE环境中使用,但有个前提:事件的回调函数必须已经被注册。你没有向/sse发送请求时,sse.js里的连接处理逻辑根本没执行,自然不会注册posted事件的回调,所以事件触发时不会有任何响应。修正方案里的全局监听是在服务器启动时就注册的,只要服务器运行就会监听,不管有没有客户端连接。

额外注意事项

  • 必须处理连接的关闭事件,及时从连接列表中移除无效连接,避免内存泄漏。
  • 加入心跳机制,防止长时间无数据时连接被反向代理或浏览器断开。
  • 使用res.flush()确保消息立即发送到客户端(部分Express版本或中间件需要)。

内容的提问来源于stack exchange,提问作者Saqlain

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 23:07:49