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

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]);

问题原因分析

  1. fetchSockets()的跨进程局限性:在集群环境中,io.fetchSockets()默认仅返回当前worker进程上的连接socket,即使使用cluster-adapter,跨worker的socket状态同步存在延迟,导致触发receiveOrder事件的worker可能无法找到连接在其他worker上的目标socket。
  2. 手动遍历查找的不可靠性:遍历所有socket的方式无法保证实时性,当订单事件触发时,目标socket的连接状态可能还未同步到当前worker,导致查找失败。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 06:37:04