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

Node.js集群实现:主端口分发请求至不同端口Worker的问题

问题解决方案

核心问题说明

  • Worker内存不共享:cluster模块的Worker是独立Node.js进程,各自拥有独立内存空间,你定义的postsDB数组会在每个Worker进程中生成副本,修改其中一个Worker的数组不会影响其他Worker。必须通过进程间通信(IPC)或外部共享存储(如Redis)实现数据统一操作。
  • 请求传递方式错误:不能直接跨进程传递req对象,序列化后会丢失socket连接等关键属性。正确流程是主进程接收请求后,将请求的核心信息(方法、路径、body等)发给Worker,Worker处理完返回响应数据,主进程再转发给客户端。

修正后的实现代码

const cluster = require("cluster");
const http = require("http");
const os = require("os");
const { parse } = require("url");

// 主进程维护共享数据库(数据量较大时建议改用Redis)
let postsDB = [{ title: "First Post" }, { title: "Second Post" }];
const WORKER_PORTS = [];
const WORKERS = [];

if (cluster.isPrimary) {
  const MAIN_PORT = 5000;
  const workerCount = os.cpus().length - 1;

  // 预创建固定数量的Worker,每个Worker监听独立端口
  for (let i = 0; i < workerCount; i++) {
    const workerPort = 5001 + i;
    WORKER_PORTS.push(workerPort);
    const worker = cluster.fork({ WORKER_PORT: workerPort });
    WORKERS.push(worker);

    // 接收Worker返回的响应,转发给对应客户端
    worker.on("message", ({ resData, clientId }) => {
      const client = clients.get(clientId);
      if (client) {
        client.writeHead(resData.statusCode, resData.headers);
        client.end(JSON.stringify(resData.body));
        clients.delete(clientId);
      }
    });

    // 处理Worker的数据库操作请求
    worker.on("message", (msg) => {
      if (msg.action === "getPosts") {
        worker.send(postsDB);
      } else if (msg.action === "addPost") {
        postsDB.push(msg.post);
        worker.send(postsDB);
      } else if (msg.action === "deletePost") {
        postsDB = postsDB.filter(p => p.title !== msg.title);
        worker.send(postsDB);
      }
    });
  }

  // 存储客户端连接,用唯一ID关联请求与响应
  const clients = new Map();
  let clientIdCounter = 0;

  // 主服务监听5000端口,分发请求
  const server = http.createServer((req, res) => {
    // 轮询选择目标Worker
    const targetWorkerIndex = clientIdCounter % WORKERS.length;
    const targetWorker = WORKERS[targetWorkerIndex];

    // 收集请求body
    let body = "";
    req.on("data", chunk => body += chunk);
    req.on("end", () => {
      const clientId = clientIdCounter++;
      clients.set(clientId, res);

      // 发送请求核心信息给Worker
      targetWorker.send({
        clientId,
        method: req.method,
        url: req.url,
        body: body ? JSON.parse(body) : null
      });
    });
  }).listen(MAIN_PORT);
  console.log(`主进程监听端口 ${MAIN_PORT}`);

  // Worker退出自动重启,保证集群稳定性
  cluster.on("exit", (worker) => {
    console.log(`Worker ${worker.process.pid} 退出,重启中...`);
    const workerIndex = WORKERS.indexOf(worker);
    const workerPort = WORKER_PORTS[workerIndex];
    const newWorker = cluster.fork({ WORKER_PORT: workerPort });
    WORKERS[workerIndex] = newWorker;

    // 重新绑定消息监听逻辑
    newWorker.on("message", ({ resData, clientId }) => {
      const client = clients.get(clientId);
      if (client) {
        client.writeHead(resData.statusCode, resData.headers);
        client.end(JSON.stringify(resData.body));
        clients.delete(clientId);
      }
    });

    newWorker.on("message", (msg) => {
      if (msg.action === "getPosts") {
        newWorker.send(postsDB);
      } else if (msg.action === "addPost") {
        postsDB.push(msg.post);
        newWorker.send(postsDB);
      } else if (msg.action === "deletePost") {
        postsDB = postsDB.filter(p => p.title !== msg.title);
        newWorker.send(postsDB);
      }
    });
  });
}

if (cluster.isWorker) {
  const WORKER_PORT = process.env.WORKER_PORT;
  // Worker启动独立端口服务(支持直接访问,也可仅处理IPC请求)
  http.createServer((req, res) => {
    handleRequest(req, res);
  }).listen(WORKER_PORT);
  console.log(`Worker ${process.pid} 监听端口 ${WORKER_PORT}`);

  // 处理主进程发来的请求
  process.on("message", ({ clientId, method, url, body }) => {
    // 模拟请求/响应对象,复用CRUD逻辑
    const mockReq = { method, url, body };
    const mockRes = {
      statusCode: 200,
      headers: { "Content-Type": "application/json" },
      body: null,
      writeHead: (code, hdrs) => {
        mockRes.statusCode = code;
        mockRes.headers = hdrs;
      },
      end: (data) => {
        mockRes.body = JSON.parse(data);
        // 将响应发回主进程
        process.send({ resData: mockRes, clientId });
      }
    };
    handleRequest(mockReq, mockRes);
  });

  // 统一CRUD处理逻辑
  function handleRequest(req, res) {
    const parsedUrl = parse(req.url);
    switch (req.method) {
      case "GET":
        // 请求主进程获取最新数据
        process.send({ action: "getPosts" });
        process.once("message", (posts) => {
          res.writeHead(200, { "Content-Type": "application/json" });
          res.end(JSON.stringify(posts));
        });
        break;
      case "POST":
        const newPost = req.body;
        // 请求主进程添加数据
        process.send({ action: "addPost", post: newPost });
        process.once("message", (updatedPosts) => {
          res.writeHead(201, { "Content-Type": "application/json" });
          res.end(JSON.stringify(updatedPosts));
        });
        break;
      case "DELETE":
        const title = parsedUrl.query.split("=")[1];
        // 请求主进程删除数据
        process.send({ action: "deletePost", title });
        process.once("message", (updatedPosts) => {
          res.writeHead(200, { "Content-Type": "application/json" });
          res.end(JSON.stringify(updatedPosts));
        });
        break;
      default:
        res.writeHead(404, { "Content-Type": "application/json" });
        res.end(JSON.stringify({ error: "Not Found" }));
    }
  }
}

关键细节说明

  1. 数据共享实现:

    • 由主进程统一维护postsDB,Worker的所有读写操作都通过IPC消息请求主进程执行,确保数据一致性。
    • 若数据量较大或需要更高性能,建议替换为Redis等专业内存数据库。
  2. 请求分发流程:

    • 主进程提前创建固定数量的Worker,每个Worker监听独立端口。
    • 主进程收到客户端请求后,用轮询算法选择Worker,发送请求核心信息。
    • Worker处理完成后将响应发回主进程,主进程再转发给对应客户端。
  3. 集群稳定性:主进程监听Worker退出事件,自动重启对应Worker,避免单个Worker故障导致服务中断。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 20:01:43