如何用Redis重构Node.js/Express的通知推送功能?
我现有一个Node.js/Express应用,通过以下代码实现了基于Server-Sent Events(SSE)的通知推送端点:
var practitionerStreams = [] // this is a list of all the streams opened by pract users to the backend async function notificationEventsHandler(req, res){ const headers ={ 'Content-Type': 'text/event-stream', 'Connection': 'keep-alive', 'Cache-Control': 'no-cache' } const practEmail = req.headers.practemail console.log("PRACT EMAIL", practEmail) const data = await ApptNotificationData.findAll({ where: { practEmail: practEmail } }) res.writeHead(200, headers) await res.write(`data:${JSON.stringify(data)} \n`) // create a new stream const newPractStream = { practEmail: practEmail, res } // add the new stream to list of streams practitionerStreams.push(newPractStream) req.on('close', () => { console.log(`${practEmail} Connection closed`); practitionerStreams = practitionerStreams.filter(pract => pract.practEmail !== pract.practEmail); }); return res } async function sendApptNotification(newNotification, practEmail){ var updatedPractitionerStream = practitionerStreams.map((stream) => // iterate through the array and find the stream that contains the pract email we want // then write the new notification to that stream { if (stream["practEmail"]==practEmail){ console.log("IF") stream.res.write(`data:${JSON.stringify(newNotification)} \n`) return stream } else { // if it doesnt contain the stream we want leave it unchanged console.log("ELSE") return stream } } ) practitionerStreams = updatedPractitionerStream }
该方案在用户连接时保留Response对象并存储,后续通过sendApptNotification推送通知,但在全量应用中性能不足。我想改用Redis重构该功能,疑问如下:
- 是否仍需保留Response对象进行写入?
- 是否要替换为前端可订阅的Redis流?
- 是否需将所有流存储在Redis中?
另外我看到示例中是将Redis事件写入Response对象,希望得到相关实现指导。
核心疑问解答
是否仍需保留Response对象进行写入?
是的,必须保留。SSE本质是HTTP长连接,前端依赖这个长连接的Response接收服务器推送的事件,所以你需要在Express请求处理函数中持有Response对象,向其写入SSE格式的数据。Redis仅作为跨进程/跨服务器的消息中间件,负责消息分发,无法替代Response的推送能力。是否要替换为前端可订阅的Redis流?
不需要。Redis Stream是服务端内部的消息队列机制,前端无法直接订阅(除非额外搭建代理,这会偏离SSE的轻量设计)。正确的做法是:前端保持SSE长连接到Express,Express服务订阅对应Redis通道,当Redis有新消息时,再通过Response推送给前端。是否需将所有流存储在Redis中?
不需要直接存储Response对象(Response是进程内资源,无法序列化到Redis)。你只需在Redis中维护「用户标识(如practEmail)」与「专属消息通道」的关联,每个Express进程仅订阅当前连接用户的通道即可。
具体实现步骤
1. 依赖准备
安装Redis客户端:
npm install redis
2. 重构SSE连接处理函数
用户连接时订阅专属Redis通道,连接关闭时取消订阅:
const { createClient } = require('redis'); // 配置Redis连接参数,比如host、port等 const redisClient = createClient({ url: 'redis://localhost:6379' }); redisClient.connect(); async function notificationEventsHandler(req, res) { const headers = { 'Content-Type': 'text/event-stream', 'Connection': 'keep-alive', 'Cache-Control': 'no-cache', 'Access-Control-Allow-Origin': '*' // 根据实际跨域需求调整 }; const practEmail = req.headers.practemail; if (!practEmail) { res.writeHead(400); return res.end('Missing practemail header'); } // 返回历史通知 const data = await ApptNotificationData.findAll({ where: { practEmail } }); res.writeHead(200, headers); // SSE格式要求消息结尾是双换行 res.write(`data:${JSON.stringify(data)}\n\n`); // 生成用户专属Redis通道名 const channel = `notifications:${practEmail}`; // 创建独立的订阅客户端(避免与发布客户端冲突) const subscriber = redisClient.duplicate(); await subscriber.connect(); // 订阅通道,收到消息时推送给前端 await subscriber.subscribe(channel, (message) => { res.write(`data:${message}\n\n`); }); // 连接关闭时清理订阅资源 req.on('close', async () => { console.log(`${practEmail} Connection closed`); await subscriber.unsubscribe(channel); await subscriber.quit(); res.end(); }); }
3. 重构通知推送函数
将新通知发布到用户专属Redis通道:
async function sendApptNotification(newNotification, practEmail) { const channel = `notifications:${practEmail}`; // 将通知序列化为JSON字符串发布到Redis通道 await redisClient.publish(channel, JSON.stringify(newNotification)); // 保留原有持久化逻辑,将通知存入数据库 await ApptNotificationData.create(newNotification); }
4. 多进程/集群场景兼容
如果应用是多进程或集群部署,每个进程会独立处理用户连接并订阅Redis通道,Redis Pub/Sub会自动将消息推送给所有订阅对应通道的进程,确保用户无论连接到哪个进程都能收到通知。
关键优化点
- 移除内存中的
practitionerStreams数组,避免内存泄漏和单进程瓶颈 - 利用Redis Pub/Sub实现跨进程消息分发,支持水平扩展
- 每个连接仅订阅自身用户的通道,减少不必要的消息接收
- 连接关闭时及时清理Redis订阅,避免资源浪费
内容的提问来源于stack exchange,提问作者Amon

