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

如何通过Node.js任务服务器优化CPU利用率?集群适配问题

Optimizing Node.js Cluster for a Variable-Workload Job Server

Great question! The basic cluster setup you shared works well for static workloads where each child process runs a fixed set of tasks, but it’s not ideal for a job server with variable task volumes—there’s no way to dynamically distribute work, so you might end up with overloaded child processes while others sit idle, wasting CPU capacity.

Here’s a better approach that turns your cluster into a master-worker task scheduler:

Core Idea

  • The master process acts as a central task manager: it holds a queue of pending jobs, monitors worker status, and distributes tasks only to idle workers.
  • Worker processes execute tasks on assignment, then notify the master once they’re done and ready for more work.
  • This way, every CPU core stays busy when there are tasks to run, and you can adapt to changing workloads seamlessly.

Implementation Example

Master Process Code

const cluster = require('cluster');
const os = require('os');
const numCPUs = os.cpus().length; // Use actual available CPU cores
const taskQueue = [];
const idleWorkers = [];

if (cluster.isMaster) {
  console.log(`Master ${process.pid} is running`);

  // Spawn initial worker pool
  for (let i = 0; i < numCPUs; i++) {
    const worker = cluster.fork();
    
    // Listen for worker status updates
    worker.on('message', (msg) => {
      if (msg.type === 'ready') {
        // Worker is idle, add to pool and assign a task if available
        idleWorkers.push(worker);
        assignPendingTask();
      } else if (msg.type === 'taskCompleted') {
        console.log(`Worker ${worker.process.pid} finished task: ${msg.taskId}`);
        // Mark worker as idle again
        idleWorkers.push(worker);
        assignPendingTask();
      }
    });

    // Restart crashed workers to maintain capacity
    worker.on('exit', (code, signal) => {
      console.log(`Worker ${worker.process.pid} crashed. Restarting...`);
      cluster.fork();
    });
  }

  // Add new tasks to the queue (replace with your actual task source: API, cron, etc.)
  function addTask(task) {
    taskQueue.push(task);
    assignPendingTask();
  }

  // Distribute pending tasks to idle workers
  function assignPendingTask() {
    while (idleWorkers.length > 0 && taskQueue.length > 0) {
      const worker = idleWorkers.shift();
      const task = taskQueue.shift();
      worker.send({ type: 'executeTask', task });
    }
  }

  // Example: Simulate variable task arrival (replace with your task trigger logic)
  setInterval(() => {
    const randomTask = {
      id: `task-${Date.now()}`,
      targetFunc: `function${Math.floor(Math.random() * 3) + 1}` // Pick function1/2/3 randomly
    };
    addTask(randomTask);
    console.log(`Added new task: ${randomTask.id}`);
  }, 1000);
}

Worker Process Code

const cluster = require('cluster');

if (cluster.isWorker) {
  console.log(`Worker ${process.pid} started`);

  // Notify master we're ready for tasks
  process.send({ type: 'ready' });

  // Listen for task assignments from master
  process.on('message', (msg) => {
    if (msg.type === 'executeTask') {
      const { task } = msg;
      // Run the assigned function
      switch (task.targetFunc) {
        case 'function1':
          runFunction1(task.id);
          break;
        case 'function2':
          runFunction2(task.id);
          break;
        case 'function3':
          runFunction3(task.id);
          break;
      }
    }
  });

  // Your actual task functions (updated to report completion)
  function runFunction1(taskId) {
    console.log(`Worker ${process.pid} executing function1 for ${taskId}`);
    // Simulate work with random delay
    setTimeout(() => {
      process.send({ type: 'taskCompleted', taskId });
    }, Math.random() * 2000);
  }

  function runFunction2(taskId) {
    console.log(`Worker ${process.pid} executing function2 for ${taskId}`);
    setTimeout(() => {
      process.send({ type: 'taskCompleted', taskId });
    }, Math.random() * 2000);
  }

  function runFunction3(taskId) {
    console.log(`Worker ${process.pid} executing function3 for ${taskId}`);
    setTimeout(() => {
      process.send({ type: 'taskCompleted', taskId });
    }, Math.random() * 2000);
  }
}

Key Optimizations for Maximum CPU Utilization

  • Dynamic Task Distribution: Tasks are only sent to idle workers, ensuring no CPU core is wasted waiting for work.
  • Auto-Scaling (Optional): Extend the master process to spawn extra workers when the task queue grows beyond a threshold, and terminate idle workers during low-traffic periods.
  • Error Resilience: Crashed workers are restarted immediately to maintain full CPU capacity.
  • Leverage All Cores: Using os.cpus().length ensures you utilize every available CPU core on your server, instead of hardcoding a fixed number.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:24:33