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

Node.js Cluster:如何等待所有工作进程的特定响应并汇总结果?

如何在Node.js Cluster中一次请求获取所有工作进程状态

这个需求很常见,核心是要把每个工作进程的响应转换成可等待的异步任务,而不是单纯靠事件监听器瞎等。我们可以用Promise来封装每个worker的回复逻辑,再结合Promise.all来等待所有结果完成,最后汇总返回给客户端。

下面直接上可运行的代码示例,我会一步步解释关键部分:

主进程代码(带Express请求处理)

const cluster = require('cluster');
const express = require('express');
const app = express();
const numCPUs = require('os').cpus().length;

if (cluster.isPrimary) {
  console.log(`主进程 ${process.pid} 正在运行`);

  // 启动所有工作进程
  for (let i = 0; i < numCPUs; i++) {
    cluster.fork();
  }

  cluster.on('exit', (worker, code, signal) => {
    console.log(`工作进程 ${worker.process.pid} 已退出`);
  });

  // 处理获取所有worker状态的请求
  app.get('/worker-stats', async (req, res) => {
    try {
      // 获取所有在线的工作进程
      const workers = Object.values(cluster.workers).filter(worker => worker.isConnected());

      if (workers.length === 0) {
        return res.json({ status: 'success', data: [], message: '无在线工作进程' });
      }

      // 给每个worker发消息,并创建对应的Promise
      const workerPromises = workers.map(worker => {
        return new Promise((resolve, reject) => {
          // 设置超时,避免某个worker挂了导致请求一直挂着
          const timeoutId = setTimeout(() => {
            reject(new Error(`工作进程 ${worker.process.pid} 未及时响应`));
          }, 5000);

          // 监听worker的回复,注意要匹配对应的worker id
          const handleMessage = (msg) => {
            if (msg.type === 'worker-status' && msg.workerId === worker.id) {
              clearTimeout(timeoutId);
              worker.off('message', handleMessage); // 移除监听器,避免内存泄漏
              resolve({
                workerId: worker.id,
                pid: worker.process.pid,
                status: msg.status
              });
            }
          };

          worker.on('message', handleMessage);
          // 向worker发送获取状态的指令
          worker.send({ type: 'get-status' });
        });
      });

      // 等待所有worker的响应
      const allStats = await Promise.all(workerPromises);
      res.json({ status: 'success', data: allStats });
    } catch (err) {
      res.status(500).json({ status: 'error', message: err.message });
    }
  });

  app.listen(3000, () => {
    console.log('服务器运行在 http://localhost:3000');
  });
} else {
  // 工作进程代码
  console.log(`工作进程 ${process.pid} 已启动`);

  // 模拟工作进程的状态数据
  let workerStatus = {
    uptime: process.uptime(),
    memoryUsage: process.memoryUsage()
  };

  // 定时更新状态(可选,模拟状态变化)
  setInterval(() => {
    workerStatus = {
      uptime: process.uptime(),
      memoryUsage: process.memoryUsage()
    };
  }, 1000);

  // 监听主进程的消息
  process.on('message', (msg) => {
    if (msg.type === 'get-status') {
      // 向主进程返回状态,带上worker id以便主进程匹配
      process.send({
        type: 'worker-status',
        workerId: cluster.worker.id,
        status: workerStatus
      });
    }
  });
}

关键逻辑解释

  1. Promise封装worker响应:
    给每个worker发送消息时,创建一个Promise。当收到该worker的对应回复时,resolve这个Promise;如果超时(比如5秒),则reject,避免请求无限挂起。

  2. 匹配worker回复:
    每个回复都带上workerId,主进程通过这个id来区分是哪个worker的回复,确保不会把A worker的回复算到B worker头上。

  3. 清理监听器:
    一旦收到对应worker的回复,就移除message事件监听器,避免内存泄漏——因为每个请求都会创建新的监听器,不用的话要及时清理。

  4. 过滤离线worker:
    用worker.isConnected()过滤掉已经断开连接的worker,避免给已经退出的worker发消息。

  5. 超时处理:
    给每个Promise设置超时时间,防止个别worker异常导致整个请求卡住。

测试方式

启动服务后,访问http://localhost:3000/worker-stats,就能看到所有在线工作进程的状态汇总结果了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 06:55:29