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

如何用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对象,希望得到相关实现指导。
重构方案与实现指导

核心疑问解答

  1. 是否仍需保留Response对象进行写入?
    是的,必须保留。SSE本质是HTTP长连接,前端依赖这个长连接的Response接收服务器推送的事件,所以你需要在Express请求处理函数中持有Response对象,向其写入SSE格式的数据。Redis仅作为跨进程/跨服务器的消息中间件,负责消息分发,无法替代Response的推送能力。

  2. 是否要替换为前端可订阅的Redis流?
    不需要。Redis Stream是服务端内部的消息队列机制,前端无法直接订阅(除非额外搭建代理,这会偏离SSE的轻量设计)。正确的做法是:前端保持SSE长连接到Express,Express服务订阅对应Redis通道,当Redis有新消息时,再通过Response推送给前端。

  3. 是否需将所有流存储在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 19:15:32