首页 / Node.js 教程 / worker_threads 工作线程

Node.js 教程

worker_threads 工作线程

本教程共 76 篇 · 第 35 篇 · 更新于 2026-07-25 · 约 7 分钟阅读

Node.jsworker_threads多线程CPU密集并发

35. worker_threads 工作线程

本节目标:用工作线程处理 CPU 密集任务,线程间通信与数据共享。

上一章讲的 cluster 是多进程方案,适合把 HTTP 请求分摊到多个核。但如果你的问题不是 “请求太多”,而是 “某个请求里的计算太重”,cluster 帮不上忙。这时候你需要的是 worker_threads——在同一个进程里开多个线程,让 CPU 密集型任务并行跑,不堵住主线程的事件循环。

线程和进程的区别

用餐厅比喻的话,进程是独立的厨房,每个厨房有自己的冰箱、灶台、厨师。线程是同一个厨房里的多个厨师,他们共用冰箱和灶台,但各自炒各自的菜。

特性worker_threadscluster (进程)
内存可共享(SharedArrayBuffer)完全隔离
启动开销高(要初始化新的 V8 实例)
通信postMessage(快)IPC(相对慢)
崩溃影响可能拖垮整个进程只影响单个 worker
适用场景CPU 密集型计算I/O 密集型服务扩展

第一个 worker 线程

worker_threads 从 v10.5.0 开始实验性引入,v12 之后稳定。v24 LTS 里可以放心用。

先写 worker 文件 heavy-task.mjs

import { parentPort, workerData } from 'node:worker_threads';

function fibonacci(n) {
  if (n <= 1) return n;
  return fibonacci(n - 1) + fibonacci(n - 2);
}

const result = fibonacci(workerData.n);
parentPort.postMessage({ input: workerData.n, result });

再写主文件 main.mjs

import { Worker } from 'node:worker_threads';
import { fileURLToPath } from 'node:url';
import { dirname, join } from 'node:path';

const __dirname = dirname(fileURLToPath(import.meta.url));

function runTask(n) {
  return new Promise((resolve, reject) => {
    const worker = new Worker(join(__dirname, 'heavy-task.mjs'), {
      workerData: { n }
    });

    worker.on('message', resolve);
    worker.on('error', reject);
    worker.on('exit', (code) => {
      if (code !== 0) reject(new Error(`worker 退出码 ${code}`));
    });
  });
}

async function main() {
  const numbers = [40, 41, 42, 43];

  // 单线程顺序跑
  console.time('单线程');
  for (const n of numbers) {
    console.log(`fib(${n}) = ${(await runTask(n)).result}`);
  }
  console.timeEnd('单线程');

  // 多线程并行跑
  console.time('多线程');
  const results = await Promise.all(numbers.map(n => runTask(n)));
  results.forEach(r => console.log(`fib(${r.input}) = ${r.result}`));
  console.timeEnd('多线程');
}

main().catch(console.error);

在我的 8 核机器上,单线程跑这四个数要 10 秒左右,多线程并行只要 3 秒。这就是 worker 的价值——把 CPU 密集型任务从主线程上卸下来。

Note

workerData 是创建 worker 时传入的初始数据,通过结构化克隆传递。parentPort.postMessage() 用来把结果发回主线程。

用同一个文件做 worker

上面的例子把 worker 逻辑拆到了单独文件里。如果你希望主线程和 worker 逻辑写在一起,可以用 isMainThread 判断:

import { Worker, isMainThread, parentPort, workerData } from 'node:worker_threads';

function computePrimes(max) {
  const sieve = new Uint8Array(max + 1);
  const primes = [];
  for (let i = 2; i <= max; i++) {
    if (!sieve[i]) {
      primes.push(i);
      for (let j = i * i; j <= max; j += i) sieve[j] = 1;
    }
  }
  return primes.length;
}

if (isMainThread) {
  const max = 5_000_000;
  const workerCount = 4;
  const chunk = Math.ceil(max / workerCount);

  const workers = [];
  for (let i = 0; i < workerCount; i++) {
    const start = i * chunk + 2;
    const end = Math.min((i + 1) * chunk + 1, max);
    workers.push(
      new Promise((resolve, reject) => {
        const worker = new Worker(import.meta.filename, {
          workerData: { start, end }
        });
        worker.on('message', resolve);
        worker.on('error', reject);
      })
    );
  }

  Promise.all(workers).then((results) => {
    const total = results.reduce((sum, r) => sum + r.count, 0);
    console.log(`1 ~ ${max} 之间的素数个数: ${total}`);
  });
} else {
  const { start, end } = workerData;
  const sieve = new Uint8Array(end + 1);
  let count = 0;
  for (let i = start; i <= end; i++) {
    if (!sieve[i]) {
      count++;
      for (let j = i * i; j <= end; j += i) sieve[j] = 1;
    }
  }
  parentPort.postMessage({ count });
}

import.meta.filename(自 v21.0+/v20.11+ 可用)指向当前文件路径,worker 可以加载自己。如果是老版本,可以用 __filenamefileURLToPath(import.meta.url)

SharedArrayBuffer:共享内存

postMessage 传递的是数据的拷贝,大对象来回传有性能损耗。如果你需要让多个线程直接读写同一块内存,可以用 SharedArrayBuffer

import { Worker } from 'node:worker_threads';

const sharedBuffer = new SharedArrayBuffer(4); // 4 字节
const counter = new Int32Array(sharedBuffer);

const workers = [];
for (let i = 0; i < 4; i++) {
  workers.push(
    new Worker(`
      const { parentPort, workerData } = require('worker_threads');
      const counter = new Int32Array(workerData.sharedBuffer);
      for (let i = 0; i < 100000; i++) {
        counter[0]++;
      }
      parentPort.postMessage('done');
    `, { eval: true, workerData: { sharedBuffer } })
  );
}

await Promise.all(workers.map(w => new Promise(r => w.on('message', r))));
console.log(`最终计数: ${counter[0]}`);

这段代码的结果大概率不是 400000。四个线程同时读、改、写 counter[0],互相踩了对方的脚。这就是经典的竞态条件

Atomics:原子操作

要安全地操作共享内存,得用 Atomics API:

import { Worker } from 'node:worker_threads';

const sharedBuffer = new SharedArrayBuffer(4);
const counter = new Int32Array(sharedBuffer);

const workers = [];
for (let i = 0; i < 4; i++) {
  workers.push(
    new Worker(`
      const { parentPort, workerData } = require('worker_threads');
      const counter = new Int32Array(workerData.sharedBuffer);
      for (let i = 0; i < 100000; i++) {
        Atomics.add(counter, 0, 1);
      }
      parentPort.postMessage('done');
    `, { eval: true, workerData: { sharedBuffer } })
  );
}

await Promise.all(workers.map(w => new Promise(r => w.on('message', r))));
console.log(`最终计数: ${Atomics.load(counter, 0)}`);

Atomics.add() 保证了对 counter[0] 的加 1 操作是原子的——读取、加 1、写回这三步不会被其他线程打断。这次结果就是准确的 400000。

Warning

SharedArrayBuffer 在 Node.js 里默认可用,但在浏览器里曾因 Spectre 漏洞被各大厂商禁用过一阵子。Node.js 服务端环境没有这个问题,但多线程编程本身就容易出 bug。不到万不得已,优先用 postMessage 传数据,代码更容易维护。

Worker Pool:复用线程

创建 worker 有开销。如果你的任务很小但数量很多,反复创建销毁 worker 反而拖慢整体速度。解决方案是维护一个 worker 池:

import { Worker } from 'node:worker_threads';
import { availableParallelism } from 'node:os';

class WorkerPool {
  constructor(workerScript, poolSize = availableParallelism()) {
    this.workerScript = workerScript;
    this.poolSize = poolSize;
    this.queue = [];
    this.workers = [];
    this.freeWorkers = [];

    for (let i = 0; i < poolSize; i++) {
      this.addWorker();
    }
  }

  addWorker() {
    const worker = new Worker(this.workerScript);
    worker.on('message', (result) => {
      if (worker.resolve) worker.resolve(result);
      worker.resolve = null;
      this.freeWorkers.push(worker);
      this.processQueue();
    });
    worker.on('error', (err) => {
      if (worker.reject) worker.reject(err);
    });
    this.workers.push(worker);
    this.freeWorkers.push(worker);
  }

  processQueue() {
    if (this.queue.length === 0 || this.freeWorkers.length === 0) return;
    const { task, resolve, reject } = this.queue.shift();
    const worker = this.freeWorkers.pop();
    worker.resolve = resolve;
    worker.reject = reject;
    worker.postMessage(task);
  }

  execute(task) {
    return new Promise((resolve, reject) => {
      this.queue.push({ task, resolve, reject });
      this.processQueue();
    });
  }

  terminate() {
    return Promise.all(this.workers.map(w => w.terminate()));
  }
}

// worker-script.mjs
// import { parentPort } from 'node:worker_threads';
// parentPort.on('message', (task) => {
//   // 做一些 CPU 密集型工作
//   const result = task.n * task.n;
//   parentPort.postMessage(result);
// });

// 使用
// const pool = new WorkerPool('./worker-script.mjs');
// const results = await Promise.all([1,2,3,4,5].map(n => pool.execute({ n })));

这个实现比较基础,生产环境可以直接用 npm 上的 workerpool 包,它处理了超时、错误重试、动态扩缩容等边界情况。

什么时候不该用 worker_threads

worker 不是万能药。以下几种情况用了反而更慢:

  1. I/O 密集型任务。文件读写、网络请求这些,Node.js 的异步 API 配合事件循环已经够高效了,再开 worker 只是增加线程切换开销。
  2. 任务本身很快。如果一次计算只要几毫秒,worker 的创建和通信开销会吃掉所有收益。
  3. 需要大量共享可变状态。一旦上了 SharedArrayBuffer,你就回到了传统多线程编程的泥潭:锁、竞态、死锁。Node.js 的单线程 + 异步模型之所以好用,就是因为它避开了这些麻烦。