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

Node集群中HTTP Upgrade请求转发至指定WebSocket实例的问题

解决Node Cluster下WebSocket请求转发问题

核心问题原因

直接用worker.send()传递request、socket这类对象时,Node会对对象进行JSON序列化,而这类对象包含大量不可序列化的底层句柄、流方法和内部状态,导致传递后丢失关键内容,无法正常完成WebSocket握手。

可行解决方案

方案1:Master作为反向代理直接转发Socket

不需要向Worker传递请求对象,Master解析请求参数后,直接建立到目标Worker的WebSocket服务的连接,将客户端Socket与目标服务的Socket双向转发,实现无缝握手。

示例代码:

// Master进程代码
const cluster = require('cluster');
const http = require('http');
const net = require('net');
const url = require('url');

if (cluster.isPrimary) {
  // 启动Worker,每个Worker对应一个WebSocket服务端口
  const workerCount = 5;
  for (let i = 1; i <= workerCount; i++) {
    const worker = cluster.fork({ PORT: 8000 + i });
    worker.port = 8000 + i;
  }

  // Master的HTTP服务器,作为单一入口
  const server = http.createServer((req, res) => {
    res.end('Use WebSocket connection');
  });

  server.on('upgrade', (req, socket, head) => {
    const query = url.parse(req.url, true).query;
    const targetPort = 8000 + parseInt(query.id);
    if (isNaN(targetPort)) {
      socket.write('HTTP/1.1 400 Bad Request\r\n\r\n');
      socket.destroy();
      return;
    }

    // 连接到目标Worker的WebSocket服务
    const proxySocket = net.connect(targetPort, 'localhost', () => {
      // 转发客户端的Upgrade请求头
      socket.write(`GET ${req.url} HTTP/1.1\r\n`);
      Object.entries(req.headers).forEach(([key, value]) => {
        socket.write(`${key}: ${value}\r\n`);
      });
      socket.write('\r\n');
      if (head.length > 0) socket.write(head);

      // 双向转发数据
      socket.pipe(proxySocket);
      proxySocket.pipe(socket);
    });

    // 处理连接错误
    proxySocket.on('error', (err) => {
      console.error(`Failed to connect to port ${targetPort}:`, err);
      socket.write('HTTP/1.1 503 Service Unavailable\r\n\r\n');
      socket.destroy();
    });

    socket.on('error', (err) => {
      console.error('Client socket error:', err);
      proxySocket.destroy();
    });
  });

  server.listen(8000, () => {
    console.log('Master HTTP server listening on port 8000');
  });
} else {
  // Worker进程代码,运行WebSocket服务
  const WebSocket = require('ws');
  const port = process.env.PORT;

  const wss = new WebSocket.Server({ port });

  wss.on('connection', (ws) => {
    console.log(`Client connected to Worker on port ${port}`);
    ws.on('message', (data) => {
      ws.send(`Received from port ${port}: ${data}`);
    });
  });

  console.log(`Worker WebSocket server running on port ${port}`);
}

方案2:通过Cluster传递Socket句柄(不序列化对象)

Node的worker.send()支持传递Socket句柄作为第二个参数,此时不会对Socket进行序列化,而是直接传递底层文件描述符,Worker收到后可以重新构建Socket并触发WebSocket升级。

示例代码:

// Master进程代码
const cluster = require('cluster');
const http = require('http');
const url = require('url');

if (cluster.isPrimary) {
  const workers = new Map();
  // 启动Worker,记录Worker与ID的映射
  for (let i = 1; i <= 5; i++) {
    const worker = cluster.fork({ WORKER_ID: i });
    workers.set(i, worker);
  }

  const server = http.createServer((req, res) => {
    res.end('Use WebSocket connection');
  });

  server.on('upgrade', (req, socket, head) => {
    const query = url.parse(req.url, true).query;
    const workerId = parseInt(query.id);
    const targetWorker = workers.get(workerId);

    if (!targetWorker) {
      socket.write('HTTP/1.1 404 Not Found\r\n\r\n');
      socket.destroy();
      return;
    }

    // 传递Upgrade事件和Socket句柄,不序列化对象
    targetWorker.send(
      { type: 'websocket-upgrade', url: req.url, headers: req.headers },
      socket // 传递Socket句柄
    );
  });

  server.listen(8000, () => {
    console.log('Master HTTP server listening on port 8000');
  });
} else {
  // Worker进程代码
  const WebSocket = require('ws');
  const http = require('http');

  const workerId = parseInt(process.env.WORKER_ID);
  // 创建一个内部HTTP服务器用于处理升级(不需要监听端口)
  const internalServer = http.createServer();
  const wss = new WebSocket.Server({ noServer: true });

  wss.on('connection', (ws) => {
    console.log(`Client connected to Worker ${workerId}`);
    ws.on('message', (data) => {
      ws.send(`Received from Worker ${workerId}: ${data}`);
    });
  });

  // 处理Master发来的Upgrade请求
  process.on('message', (msg, socket) => {
    if (msg.type !== 'websocket-upgrade' || !socket) return;

    // 手动触发内部HTTP服务器的upgrade事件
    internalServer.emit('upgrade', {
      url: msg.url,
      headers: msg.headers,
      method: 'GET'
    }, socket, Buffer.alloc(0));
  });

  // 绑定WebSocket升级逻辑到内部服务器
  internalServer.on('upgrade', (req, socket, head) => {
    wss.handleUpgrade(req, socket, head, (ws) => {
      wss.emit('connection', ws, req);
    });
  });

  console.log(`Worker ${workerId} ready to handle WebSocket connections`);
}

注意事项

  • 要处理端口不存在、Worker离线等异常情况,避免客户端Socket挂起或内存泄漏。
  • 方案1的反向代理方式更直观,不需要Worker做额外的事件监听,适合大多数场景。
  • 方案2利用Cluster的句柄传递特性,减少了一次Socket转发,但需要Worker配合处理消息事件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 00:55:04