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

Node.js中如何用多Worker处理HTTP请求任务?Cluster与PM2选型疑问

Node.js 任务并行执行方案选择与实现

一、方案选择:Cluster vs PM2

  • 原生Cluster模块:适合代码层面精准控制进程调度,比如根据任务类型动态调整工作进程数,无需额外依赖,适配自定义需求较强的场景。
  • PM2:封装了Cluster的进程管理工具,提供开箱即用的进程守护、负载均衡、日志管理等功能,适合快速部署,不想自己写进程管理逻辑的场景。

你的任务是周期性从数据库取动态数据并并行处理,两种方案都能实现,下面分别给出具体实现方式。

二、原生Cluster实现步骤

1. 完整代码实现

主进程负责数据库查询和任务分发,工作进程负责具体任务处理:

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

if (cluster.isPrimary) {
  // 启动与CPU核心数匹配的工作进程
  for (let i = 0; i < numCPUs; i++) {
    cluster.fork();
  }

  // 周期性分发任务
  async function dispatchTasks() {
    while (true) {
      const requestIdArray = await db.find();
      // 按工作进程数拆分任务分片
      const taskChunks = [];
      const chunkSize = Math.ceil(requestIdArray.length / numCPUs);
      for (let i = 0; i < requestIdArray.length; i += chunkSize) {
        taskChunks.push(requestIdArray.slice(i, i + chunkSize));
      }

      // 等待所有工作进程完成当前批次任务
      await Promise.all(
        Object.values(cluster.workers).map((worker, index) => {
          return new Promise((resolve) => {
            if (!taskChunks[index]) return resolve();
            worker.send({ type: 'TASK', data: taskChunks[index] });
            worker.once('message', (msg) => {
              if (msg.type === 'TASK_DONE') resolve();
            });
          });
        })
      );

      // 避免循环过于频繁,可根据需求调整间隔
      await new Promise(resolve => setTimeout(resolve, 1000));
    }
  }

  dispatchTasks().catch(err => console.error('任务分发失败:', err));

  // 工作进程崩溃自动重启
  cluster.on('exit', (worker) => {
    console.log(`工作进程 ${worker.process.pid} 退出,重启中...`);
    cluster.fork();
  });
} else {
  // 工作进程处理任务逻辑
  process.on('message', async (msg) => {
    if (msg.type === 'TASK') {
      try {
        // 把分片任务再拆成每10个一组,并行处理
        const workingIdChunks = [];
        for (let i = 0; i < msg.data.length; i += 10) {
          workingIdChunks.push(msg.data.slice(i, i + 10));
        }
        await Promise.all(
          workingIdChunks.map(async (chunk) => {
            const apiRes = await thirdPartyAPIRequest(chunk);
            await work(apiRes);
          })
        );
        // 通知主进程任务完成
        process.send({ type: 'TASK_DONE' });
      } catch (err) {
        console.error(`工作进程 ${process.pid} 处理任务失败:`, err);
        process.send({ type: 'TASK_DONE' }); // 即使失败也通知,避免阻塞主进程
      }
    }
  });
}

2. 关键说明

  • 主进程统一查询数据库,避免多进程重复查询造成资源浪费
  • 工作进程内用Promise.all并行处理小批次任务,最大化利用异步IO能力
  • 自带进程崩溃重启机制,保证任务稳定性

三、PM2实现步骤

PM2无需手动写进程管理逻辑,只需调整业务代码实现并行,再通过PM2启动多进程即可。

1. 业务代码修改

async function processTasks() {
  while (true) {
    const requestIdArray = await db.find();
    // 拆分为每10个ID一组,并行处理
    const workingIdChunks = [];
    for (let i = 0; i < requestIdArray.length; i += 10) {
      workingIdChunks.push(requestIdArray.slice(i, i + 10));
    }
    await Promise.all(
      workingIdChunks.map(async (chunk) => {
        const apiRes = await thirdPartyAPIRequest(chunk);
        await work(apiRes);
      })
    );
    // 任务间隔,按需调整
    await new Promise(resolve => setTimeout(resolve, 1000));
  }
}

processTasks().catch(err => console.error('任务处理失败:', err));

2. PM2启动命令

安装PM2后,执行以下命令启动多进程:

pm2 start your-script.js -i max
  • -i max表示根据CPU核心数启动对应数量的进程
  • 注意:多进程会独立执行查询逻辑,若想避免重复查询,可单独用一个进程做DB查询并将任务推送到Redis队列,其他进程从队列取任务处理。

四、进阶优化:任务队列

如果任务量较大或需要分布式调度,推荐结合Redis任务队列(如bullmq):

  1. 单独的调度进程查询数据库,将任务推送到队列
  2. 多工作进程(Cluster/PM2启动)从队列取任务并行处理
  3. 这种方式能彻底避免重复查询,还能灵活控制并发数,扩展性更强

方案对比

方案优势劣势适用场景
原生Cluster代码可控性高,无额外依赖需要自行编写进程管理、任务分发逻辑自定义需求强,需精准控制调度
PM2开箱即用,自带进程守护、日志管理任务调度逻辑依赖PM2,自定义程度低快速部署,无需复杂进程管理
任务队列+多进程分布式处理,扩展性强,避免重复查询DB需要额外依赖(Redis等),增加复杂度任务量大,需分布式调度的场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 21:42:44