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发送请求,却想执行其中的函数。
问题分析与解决方案
问题原因
- 回调未被注册:只有当客户端发起
/sse请求时,sse.js里的PostEvent.on('posted', ...)才会执行绑定。如果没有任何客户端连接到/sse,这个回调根本不存在,事件触发时自然没有响应。 - 重复绑定隐患:每次新的
/sse连接进来,都会重复绑定posted事件,导致后续一个事件触发时,多个相同回调执行,还可能引发内存泄漏。
修正方案:维护活跃SSE连接列表
正确的做法是把所有活跃的SSE连接统一管理,当posted事件触发时,遍历所有连接推送消息,而不是在每个连接里单独绑定事件。
修正后的代码:
- postEvent.js(不变)
const EventEmitter = require("node:events"); class PostEvent extends EventEmitter {} module.exports = new PostEvent();
- 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?.(); }); });
- 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
相关产品推荐
相关产品推荐

