Redis订阅事件触发后如何让Socket.IO客户端正确加入指定房间
核心问题说明
你当前遇到两个核心问题:
- 全局变量
room、socketId属于所有连接共享的状态,多客户端连接时会互相覆盖,导致状态异常 - Redis订阅事件和Socket连接事件没有做关联绑定,时序不匹配导致加入房间逻辑触发时机错误
你需要保证Redis发布的conversation消息中携带用于识别目标用户的唯一标识(比如userId),Socket客户端连接时需要将该标识通过auth参数或其他方式传递到服务端,才能完成消息和目标Socket的匹配。
方案1:单实例部署使用Node原生EventEmitter实现
该方案无需额外依赖,通过自定义事件关联Redis消息和Socket操作,适合单节点部署的场景:
const express = require('express'); const app = express(); const socketio = require('socket.io'); const redis = require("redis"); const { EventEmitter } = require('events'); // 初始化事件总线 const eventBus = new EventEmitter(); const expressServer = app.listen(3001, () => console.log("Express is running!")); const io = socketio(expressServer); const sub = redis.createClient({port: 6379, host: '127.0.0.1'}); sub.subscribe('conversation'); sub.on('error', function (error) { console.log('ERROR ' + error); }); sub.on('connect', function(){ console.log('Redis client connected'); }); // Redis订阅消息回调 sub.on('message', async function (channel, data) { data = JSON.parse(data); if(channel === 'conversation'){ const roomId = data.id; const targetUserId = data.userId; // 消息携带要加入房间的用户ID // 触发自定义加入房间事件 eventBus.emit('join_room', roomId, targetUserId); } }); io.on('connect', (socket) => { // 客户端连接时携带当前用户的userId,和Redis消息中的标识匹配 const userId = socket.handshake.auth.userId; // 监听加入房间事件 const joinRoomHandler = (roomId, targetUserId) => { if (targetUserId === userId) { socket.join(roomId); } }; eventBus.on('join_room', joinRoomHandler); // Socket断开时移除事件监听,避免内存泄漏 socket.on('disconnect', () => { eventBus.off('join_room', joinRoomHandler); }); }); // 房间加入成功回调全局只需要注册一次,不要放在connect回调里 io.of("/").adapter.on("join-room", (room, id) => { console.log(`socket ${id} has joined room ${room}`); });
方案2:多实例部署使用Redis适配器remoteJoin实现
该方案适配多节点集群部署场景,通过remoteJoin方法可以跨实例操作Socket加入房间,解决你之前全局socketId被覆盖的问题:
const express = require('express'); const app = express(); const socketio = require('socket.io'); const redis = require("redis"); const redisAdapter = require('@socket.io/redis-adapter'); const expressServer = app.listen(3001, () => console.log("Express is running!")); const io = socketio(expressServer); const sub = redis.createClient({port: 6379, host: '127.0.0.1'}); const pub = sub.duplicate(); io.adapter(redisAdapter(pub, sub)); sub.subscribe('conversation', function(){ console.log('Redis subscribed success') }) // 存储userId和socketId的映射,代替全局socketId const userSocketMap = new Map(); sub.on('error', function (error) { console.log('ERROR ' + error); }); sub.on('connect', function(){ console.log('Redis client connected'); }); sub.on('message', async function (channel, data) { data = JSON.parse(data); if(channel === 'conversation'){ const roomId = data.id; const targetUserId = data.userId; // 从映射表中获取对应用户的socketId const targetSocketId = userSocketMap.get(targetUserId); if (targetSocketId) { try { // remoteJoin是异步方法,必须加await await io.of('/').adapter.remoteJoin(targetSocketId, roomId); console.log(`socket ${targetSocketId} joined room ${roomId}`); } catch (e) { // 捕获异常,处理socket已断开等异常场景 console.error('join room failed:', e) } } } }); // 全局只需要注册一次join-room回调,不要放在connect回调里 io.of("/").adapter.on("join-room", (room, id) => { console.log(`socket ${id} has joined room ${room}`); }); io.on('connect', (socket) => { const userId = socket.handshake.auth.userId; // 建立用户和socketId的映射 userSocketMap.set(userId, socket.id); // Socket断开时删除映射,避免内存泄漏 socket.on('disconnect', () => { userSocketMap.delete(userId); }); });
注意事项
remoteJoin为异步方法,必须使用await调用,否则会出现时序异常join-room事件回调全局仅需注册一次,不要放在Socket的connect回调中,避免重复注册
内容的提问来源于stack exchange,提问作者Kyle Corbin Hurst
相关产品推荐
相关产品推荐

