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

如何在Express+ws集群环境中向不同Worker的所有客户端广播消息?

跨Worker的WebSocket广播解决方案

问题场景

在Express应用中使用cluster模块和ws包搭建WebSocket服务时,客户端连接会被分配到不同的Worker进程处理。当前实现的广播功能只能将消息发送给同一Worker进程下的客户端,无法覆盖其他Worker上的连接,而单进程模式下功能正常。

核心原因

cluster模式下每个Worker都是独立的Node.js进程,各自维护自己的wss.clients客户端集合,进程间内存完全隔离,所以单个Worker的广播无法触达其他Worker的客户端。

实现方案

通过Master进程作为消息中转站,基于cluster的IPC(进程间通信)实现跨Worker广播:

  • 当某个Worker收到客户端的广播消息时,先将消息发送给Master进程
  • Master进程收到消息后,转发给所有在线的Worker进程
  • 每个Worker收到Master的转发指令后,向自己管理的所有客户端广播消息

修改后的完整代码

主运行函数

const run = () => {
  if (cluster.isPrimary) {
    // Master进程:创建Worker并处理消息转发
    for (let index = 0; index < numberCPUs; index++) {
      cluster.fork();
    }

    // 接收Worker发来的广播请求,转发给所有其他Worker
    cluster.on('message', (worker, msg) => {
      if (msg.type === 'broadcast') {
        Object.values(cluster.workers).forEach(targetWorker => {
          if (targetWorker.id !== worker.id) {
            targetWorker.send({
              type: 'broadcast',
              data: msg.data
            });
          }
        });
      }
    });

    cluster.on("exit", (worker, code, signal) => {
      console.error(`worker ${worker.process.pid} died (${signal || code}). restarting it in a sec`);
      setTimeout(() => cluster.fork(), 1000);
    });
  } else {
    console.log(`Worker ${cluster.worker.id} : started`);
    const http = require("node:http");
    const WebSocket = require("ws");

    const server = http.createServer(application);
    const wss = new WebSocket.Server({ server });

    // 监听Master发来的广播指令,执行本地客户端广播
    process.on('message', (msg) => {
      if (msg.type === 'broadcast') {
        broadcast(WebSocket, wss, msg.data);
      }
    });

    wss.on("connection", (ws, request, client) => {
      console.log(`Worker ${cluster.worker.id} handled new connection !`);
      ws.on("message", (message) => {
        console.log(`Worker ${cluster.worker.id} : Received message ${message}`);
        // 先给当前Worker下的客户端广播
        broadcast(WebSocket, wss, message);
        // 将消息发给Master,让其转发给其他Worker
        process.send({
          type: 'broadcast',
          data: message.toString()
        });
      });
    });

    server.listen(port, () => {
      console.log(`Server is running on ${port} ...`);
    });
  }
};

广播函数(无需修改)

const broadcast = (webSocket, wss, message) => {
  for (const client of wss.clients) {
    if (client.readyState === webSocket.OPEN) {
      client.send(message, { binary: false });
    }
  }
};

方案说明

  1. Master进程:作为消息枢纽,接收任意Worker的广播请求后,转发给所有其他Worker,避免重复向发送消息的Worker转发。
  2. Worker进程:
    • 收到客户端消息时,先给本地客户端广播一次(减少延迟)
    • 同时将消息发送给Master,由Master完成跨Worker转发
    • 监听Master的广播指令,收到后向自己管理的客户端同步消息
  3. 注意将message转为字符串传输,避免IPC处理二进制数据可能出现的异常(若需二进制传输,可调整为Buffer格式)

内容的提问来源于stack exchange,提问作者Alirezakvr

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 17:17:31