如何基于WebSocket与Redis实现类Instagram的发布订阅功能
用WebSocket + Redis实现粉丝动态实时通知功能
核心逻辑概述
用户发布内容后要让粉丝实时收到通知,核心是靠WebSocket维持客户端与服务端的长连接,用Redis做消息路由和状态管理——Redis的Pub/Sub能高效实现消息广播,同时它的内存结构可以方便存储在线用户状态、粉丝列表和离线消息。
1. 维护用户在线状态与连接映射
用户登录并建立WebSocket连接时,服务端要做两件事:
- 把用户ID存入Redis的Set集合
online_users,标记用户在线;如果支持多端登录,Set会自动保留唯一用户ID(若要追踪多连接,也可以用Hash,key是用户ID,value是连接ID的集合)。 - 为该用户订阅Redis中专属的通知频道
user_notifications:{用户ID},后续给这个用户发的通知都会推送到这个频道。
2. 粉丝列表的缓存优化
提前把每个用户的粉丝ID存在Redis的Set集合followers:{用户ID}里,比如用户B关注用户A时,执行SADD followers:A B。这样用户A发布内容时,不用每次查数据库,直接从Redis取粉丝列表,速度更快。
3. 发布内容时的消息推送流程
当用户A发布一条新内容:
- 先把内容存到数据库,生成唯一的动态ID;
- 从Redis的
followers:A取出所有粉丝ID; - 遍历每个粉丝ID:
- 用
SISMEMBER online_users {粉丝ID}判断粉丝是否在线; - 如果在线,构造通知消息(比如动态摘要、作者ID、发布时间),调用Redis的
PUBLISH user_notifications:{粉丝ID} 消息内容,把消息推到粉丝的专属频道; - 如果不在线,把消息存入Redis的有序集合
offline_notifications:{粉丝ID},按时间戳排序,等用户下次上线再推送。
- 用
4. WebSocket服务端的消息转发
WebSocket服务端会监听自己负责的用户的专属频道,一旦Redis频道传来消息,就立刻通过对应的WebSocket连接把消息推给前端。前端收到消息后,直接更新动态流UI(比如在顶部插入新动态,或者弹出小红点提示)。
5. 关键细节处理
- 多服务节点兼容:如果部署多个WebSocket服务实例,Redis Pub/Sub会把消息广播给所有节点,每个节点只需要把消息推送给自己管理的在线用户,不会漏推或重复推。
- 连接断开清理:用户关闭页面或断开WebSocket时,服务端要从
online_users移除该用户ID,同时取消订阅对应的通知频道,避免无效资源占用。 - 离线消息同步:用户上线时,服务端从
offline_notifications:{用户ID}取出所有未读消息,推送给用户后清空该集合(或者标记已读,根据产品需求)。 - 消息去重:给每条通知加唯一ID,前端收到后记录已处理的ID,防止网络波动导致的重复推送。
简单代码示例(Node.js)
初始化Redis与WebSocket服务
const Redis = require('ioredis'); const WebSocket = require('ws'); // Redis客户端 const redis = new Redis(); // WebSocket服务 const wss = new WebSocket.Server({ port: 8080 });
处理WebSocket连接
wss.on('connection', async (ws, req) => { // 从请求中解析用户ID(比如通过token验证) const userId = parseUserIdFromRequest(req); if (!userId) { ws.close(); return; } // 标记用户在线 await redis.sadd('online_users', userId); // 订阅用户专属通知频道 const subscriber = new Redis(); const channel = `user_notifications:${userId}`; subscriber.subscribe(channel); // 收到频道消息时推送给前端 subscriber.on('message', (ch, msg) => { if (ch === channel && ws.readyState === WebSocket.OPEN) { ws.send(msg); } }); // 连接断开时清理资源 ws.on('close', async () => { await redis.srem('online_users', userId); subscriber.unsubscribe(channel); subscriber.quit(); }); });
处理内容发布与通知推送
async function publishNewPost(userId, postContent) { // 1. 保存内容到数据库 const postId = await savePostToDatabase(userId, postContent); // 2. 获取粉丝列表 const followers = await redis.smembers(`followers:${userId}`); // 3. 构造通知消息 const notification = JSON.stringify({ type: 'new_post', postId, authorId: userId, content: postContent.slice(0, 50) + '...', timestamp: Date.now() }); // 遍历粉丝推送消息 for (const fanId of followers) { const isOnline = await redis.sismember('online_users', fanId); if (isOnline) { await redis.publish(`user_notifications:${fanId}`, notification); } else { // 离线消息存入有序集合 await redis.zadd(`offline_notifications:${fanId}`, Date.now(), notification); } } }
内容的提问来源于stack exchange,提问作者kpatel23
相关产品推荐
相关产品推荐

