worker_threads 工作线程
本教程共 76 篇 · 第 35 篇 · 更新于 2026-07-25 · 约 7 分钟阅读
35. worker_threads 工作线程
本节目标:用工作线程处理 CPU 密集任务,线程间通信与数据共享。
上一章讲的 cluster 是多进程方案,适合把 HTTP 请求分摊到多个核。但如果你的问题不是 “请求太多”,而是 “某个请求里的计算太重”,cluster 帮不上忙。这时候你需要的是 worker_threads——在同一个进程里开多个线程,让 CPU 密集型任务并行跑,不堵住主线程的事件循环。
线程和进程的区别
用餐厅比喻的话,进程是独立的厨房,每个厨房有自己的冰箱、灶台、厨师。线程是同一个厨房里的多个厨师,他们共用冰箱和灶台,但各自炒各自的菜。
| 特性 | worker_threads | cluster (进程) |
|---|---|---|
| 内存 | 可共享(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 可以加载自己。如果是老版本,可以用 __filename 或 fileURLToPath(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 不是万能药。以下几种情况用了反而更慢:
- I/O 密集型任务。文件读写、网络请求这些,Node.js 的异步 API 配合事件循环已经够高效了,再开 worker 只是增加线程切换开销。
- 任务本身很快。如果一次计算只要几毫秒,worker 的创建和通信开销会吃掉所有收益。
- 需要大量共享可变状态。一旦上了
SharedArrayBuffer,你就回到了传统多线程编程的泥潭:锁、竞态、死锁。Node.js 的单线程 + 异步模型之所以好用,就是因为它避开了这些麻烦。