Socket.IO向客户端发送消息行为不一致问题排查求助
Socket.IO集群环境下消息丢失问题排查与解决
场景与当前流程
- 用户下单后向服务器触发
receiveOrder事件 - 服务器向对应餐厅客户端触发
newOrder事件,餐厅客户端接收新订单数据
问题现象
向餐厅客户端发送订单数据时,时而能正常接收,时而出现消息丢失的情况。目前通过连接时附带的user-ID查找目标客户端socket,此前用数组存储socket存在跨进程同步问题,改用const allSocketUsers = await io.fetchSockets();遍历查找。
环境说明
- Express服务器运行在集群环境,已集成
@socket.io/cluster-adapter和@socket.io/sticky处理集群调度 - 客户端为React Native应用,连接时通过query参数携带
user-ID和role信息
服务器端配置代码
const numCPUs = os.cpus().length; const PORT = process.env.PORT || 3000; env.config(); if (cluster.isPrimary) { console.log(`Primary ${process.pid} is running`); // Setup Socket.IO primary adapter setupPrimary(); const httpServer = createServer(); setupMaster(httpServer, { loadBalancingMethod: "least-connection", }); httpServer.listen(PORT, () => { console.log(`Server is listening on port ${PORT}`); }); // Fork workers for (let i = 0; i < numCPUs; i++) { cluster.fork(); } cluster.on("exit", (worker, code, signal) => { console.log(`Worker ${worker.process.pid} died. Forking a new one.`); cluster.fork(); }); } else { const app = express(); app.use(morgan("dev")); app.use(express.json()); app.use(express.urlencoded({ extended: true })); app.use(cookieParser()); app.use( cors({ origin: "*", methods: ["GET", "POST", "PUT", "DELETE"], }) ); // Setup routes app.use("/api/user", UserRoutes); // rest of app routes...... /* create server */ const httpServer = createServer(app); const io = new Server(httpServer, { cors: { origin: "*", methods: ["GET", "POST", "PUT", "DELETE"], }, }); /* attach socket io cluster to handle multi threading */ io.adapter(createAdapter()); // Setup sticky session worker Thus client will always connect to same worker server setupWorker(io); const PORTWorker = parseInt(process.env.PORT) || 3000 + cluster.worker.id; httpServer.listen(PORTWorker, () => { console.log(`Worker ${process.pid} started on port ${PORTWorker}`); }); /* ====================START============== */ /* socket implmenetation */ io.on("connection", async (socket: Socket) => { const id = socket?.handshake?.query["user-ID"] as string; // add userId to connected socket socket.data = { userId: id.toString(), }; socket.on("receiveOrder", async (data) => { const { clientId, order } = data; const allConnectedSockets = await io.fetchSockets(); const restaurantSocket = allConnectedSockets?.find((userSocket) => { return userSocket.data.userId == clientId; }); if (restaurantSocket) { //emit notifyClient to the client restaurantSocket.emit("newOrder", { ...order }); } }); socket.on("disconnect", () => { socket.disconnect(); console.log( "client disconnected: ", socket?.handshake?.query["user-ID"] ); }); }); /* ======================[END WS]===================== */ }
客户端(餐厅)监听代码
useEffect(() => { const socket = io.connect('https://xxxx.xx', { query: { 'user-ID': user.user_id, role: 'restaurant', }, }); socket.on('connect', () => { console.log('restaurant socket connected'); }); //listening to incoming order socket.on('newOrder', async (data) => { console.log('new order created', data); //handle incoming new orders ... }); return () => { socket.disconnect(); }; }, [isFoucsed]);
问题原因分析
fetchSockets()的跨进程局限性:在集群环境中,io.fetchSockets()默认仅返回当前worker进程上的连接socket,即使使用cluster-adapter,跨worker的socket状态同步存在延迟,导致触发receiveOrder事件的worker可能无法找到连接在其他worker上的目标socket。- 手动遍历查找的不可靠性:遍历所有socket的方式无法保证实时性,当订单事件触发时,目标socket的连接状态可能还未同步到当前worker,导致查找失败。
- sticky session潜在问题:若
@socket.io/sticky配置存在偏差,可能导致同一客户端连接到不同worker,进一步增加跨进程查找socket的失败概率。
解决方案
方案1:使用Socket.IO房间(Room)机制(推荐)
每个餐厅客户端连接时,加入以自身user-ID命名的专属房间,发送消息时直接向房间推送,无需手动查找socket:
io.on("connection", async (socket: Socket) => { const id = socket?.handshake?.query["user-ID"] as string; socket.data = { userId: id.toString(), }; // 加入以user-ID命名的房间 socket.join(id.toString()); socket.on("receiveOrder", async (data) => { const { clientId, order } = data; // 直接向目标房间发送事件,集群环境下adapter会自动同步到对应worker io.to(clientId).emit("newOrder", { ...order }); }); socket.on("disconnect", () => { console.log("client disconnected: ", socket?.handshake?.query["user-ID"]); }); });
方案2:分布式socket映射存储
使用Redis等分布式存储维护user-ID与socket.id的映射关系,确保跨worker进程能获取到目标socket的信息:
- 连接时:将
socket.id和user-ID的键值对存入Redis - 发送时:通过
user-ID从Redis查询对应的socket.id,再调用io.to(socketId).emit(...)发送消息 - 断开时:从Redis删除对应映射,避免脏数据
方案3:校验sticky session配置
确认@socket.io/sticky的setupWorker配置正确,确保同一客户端始终连接到同一个worker,减少跨进程通信的复杂度,提升本地查找socket的成功率。
内容的提问来源于stack exchange,提问作者Alfe
相关产品推荐
相关产品推荐

