Node.js: Worker 线程
最后更新:2026-08-26
1. 你将学到
- 使用
Worker类创建工作线程并管理生命周期 - 通过
parentPort实现主线程与 Worker 双向消息通信 - 利用
workerData向 Worker 传递初始化数据 - 使用
MessageChannel/MessagePort建立线程间直接通信 - 通过
SharedArrayBuffer与Atomics实现共享内存操作 - 构建线程池模式复用 Worker 提升吞吐量
- 对比 Worker Threads、child_process 与 cluster 的适用场景
2. 故事:Alice 的图像处理加速
Alice 在一家 SaaS 公司负责图像处理服务。用户上传图片后,系统需要为每张图片生成 3 种尺寸的缩略图。最初的单线程方案处理 100 张图片需要约 30 秒,高峰期请求排队严重。
她决定用 Worker Threads 将任务分配到 4 个线程并行处理:主线程读取文件列表,通过 workerData 传入任务参数,Worker 完成缩略图生成后通过 parentPort 返回结果。最终 100 张图片的处理时间降到约 8 秒,吞吐量提升近 4 倍。
3. Worker 类基础
worker_threads 模块是 Node.js 提供的真正的多线程能力。每个 Worker 在独立线程中执行脚本,拥有自己的 V8 引擎实例和事件循环。
const { Worker } = require('worker_threads');
const worker = new Worker('./heavy-task.js', {
workerData: { taskId: 42, input: 'hello' }
});
worker.on('message', (result) => {
console.log('Result from worker:', result);
});
worker.on('error', (err) => {
console.error('Worker error:', err);
});
worker.on('exit', (code) => {
if (code !== 0) {
console.error(`Worker stopped with exit code ${code}`);
}
});
Worker 脚本内部:
// heavy-task.js
const { parentPort, workerData } = require('worker_threads');
const { taskId, input } = workerData;
const result = performHeavyComputation(input);
parentPort.postMessage({ taskId, result });
function performHeavyComputation(data) {
let hash = 0;
for (let i = 0; i < 1e8; i++) {
hash = (hash + data.charCodeAt(i % data.length)) % 65536;
}
return hash;
}
▶ 示例:最简 Worker 创建与通信
// main.js
const { Worker } = require('worker_threads');
const worker = new Worker(`
const { parentPort } = require('worker_threads');
parentPort.postMessage('Hello from worker!');
`, { eval: true });
worker.on('message', (msg) => {
console.log(msg); // Hello from worker!
});
4. parentPort 双向通信
parentPort 是 Worker 与主线程之间的消息通道。主线程通过 worker.postMessage() 发送,Worker 通过 parentPort.postMessage() 回复。双方都监听 message 事件接收数据。
▶ 示例:双向消息交互
// main.js
const { Worker } = require('worker_threads');
const worker = new Worker('./echo-worker.js');
worker.postMessage({ action: 'greet', name: 'Alice' });
worker.on('message', (msg) => {
console.log('Main received:', msg);
if (msg.action === 'greet_reply') {
worker.postMessage({ action: 'task', payload: 100 });
}
});
// echo-worker.js
const { parentPort } = require('worker_threads');
parentPort.on('message', (msg) => {
if (msg.action === 'greet') {
parentPort.postMessage({
action: 'greet_reply',
text: `Hello, ${msg.name}!`
});
} else if (msg.action === 'task') {
const result = msg.payload * 2;
parentPort.postMessage({ action: 'result', value: result });
}
});
▶ 示例:Transferable 对象零拷贝传输
const { Worker } = require('worker_threads');
const buffer = new ArrayBuffer(1024 * 1024); // 1 MB
const worker = new Worker('./process-buffer.js');
worker.postMessage({ buffer }, [buffer]);
console.log('Buffer transferred, main thread no longer owns it');
// process-buffer.js
const { parentPort, workerData } = require('worker_threads');
parentPort.on('message', ({ buffer }) => {
const view = new Uint8Array(buffer);
view[0] = 42;
parentPort.postMessage({ done: true, firstByte: view[0] }, [buffer]);
});
5. workerData 初始化数据
workerData 是在创建 Worker 时通过选项传入的只读初始数据。Worker 内部直接读取,无需消息传递。适合传递配置、文件路径、任务参数等。
▶ 示例:批量图片处理 Worker
// main.js
const { Worker } = require('worker_threads');
const path = require('path');
const imageFiles = [
'photo-001.jpg', 'photo-002.jpg', 'photo-003.jpg',
'photo-004.jpg', 'photo-005.jpg', 'photo-006.jpg'
];
const WORKER_COUNT = 4;
const chunkSize = Math.ceil(imageFiles.length / WORKER_COUNT);
for (let i = 0; i < WORKER_COUNT; i++) {
const chunk = imageFiles.slice(i * chunkSize, (i + 1) * chunkSize);
const worker = new Worker('./image-worker.js', {
workerData: {
workerId: i,
files: chunk,
outputDir: './thumbnails'
}
});
worker.on('message', (result) => {
console.log(`Worker ${i} done:`, result.processed);
});
}
// image-worker.js
const { parentPort, workerData } = require('worker_threads');
const path = require('path');
const { workerId, files, outputDir } = workerData;
async function generateThumbnail(file) {
// Simulate image processing
return new Promise((resolve) => {
setTimeout(() => resolve(`${file} -> thumb`), 100);
});
}
(async () => {
const processed = [];
for (const file of files) {
const result = await generateThumbnail(file);
processed.push(result);
}
parentPort.postMessage({ workerId, processed });
})();
6. MessageChannel 与 MessagePort
MessageChannel 创建一对互通的 MessagePort,可分配给不同 Worker,实现线程间直接通信,无需经过主线程中转。
graph LR
Main["Main Thread"] -->|"worker.postMessage()"| W1["Worker 1"]
W1 -->|"parentPort.postMessage()"| Main
Main -->|"worker.postMessage()"| W2["Worker 2"]
W2 -->|"parentPort.postMessage()"| Main
W1 <-->|"MessagePort"| W2
▶ 示例:两个 Worker 直接通信
// main.js
const { Worker, MessageChannel } = require('worker_threads');
const worker1 = new Worker('./chat-worker.js', {
workerData: { id: 1 }
});
const worker2 = new Worker('./chat-worker.js', {
workerData: { id: 2 }
});
const { port1, port2 } = new MessageChannel();
worker1.postMessage({ port: port1 }, [port1]);
worker2.postMessage({ port: port2 }, [port2]);
// chat-worker.js
const { parentPort, workerData } = require('worker_threads');
parentPort.once('message', ({ port }) => {
port.on('message', (msg) => {
console.log(`Worker ${workerData.id} received:`, msg);
});
setInterval(() => {
port.postMessage(`Hello from Worker ${workerData.id}`);
}, 1000);
});
7. SharedArrayBuffer 与 Atomics
SharedArrayBuffer 允许多个线程共享同一块内存。配合 Atomics 提供的原子操作,可安全地读写共享数据,避免竞态条件。
▶ 示例:共享计数器
// main.js
const { Worker } = require('worker_threads');
const sharedBuffer = new SharedArrayBuffer(4);
const sharedArray = new Int32Array(sharedBuffer);
const WORKER_COUNT = 4;
const INCREMENTS_PER_WORKER = 100000;
for (let i = 0; i < WORKER_COUNT; i++) {
const worker = new Worker('./counter-worker.js', {
workerData: { sharedBuffer, increments: INCREMENTS_PER_WORKER }
});
worker.on('exit', () => {
const final = Atomics.load(sharedArray, 0);
console.log(`Final counter value: ${final}`);
});
}
// counter-worker.js
const { parentPort, workerData } = require('worker_threads');
const { sharedBuffer, increments } = workerData;
const sharedArray = new Int32Array(sharedBuffer);
for (let i = 0; i < increments; i++) {
Atomics.add(sharedArray, 0, 1);
}
parentPort.postMessage('done');
8. 通信方式对比
| 通信方式 | 方向 | 特点 | 适用场景 |
|---|---|---|---|
parentPort |
主线程 ↔ Worker | 双向消息,结构化克隆 | 一般任务通信 |
workerData |
主线程 → Worker | 只读,创建时传入 | 初始化配置/参数 |
MessageChannel |
Worker ↔ Worker | 直连端口,不经主线程 | 线程间协作 |
SharedArrayBuffer |
所有线程共享 | 零拷贝,需 Atomics | 高频数据交换 |
9. Worker vs child_process vs cluster
| 特性 | Worker Threads | child_process | cluster |
|---|---|---|---|
| 单位 | 线程 | 进程 | 进程 |
| 内存 | 共享进程内存 | 独立内存空间 | 独立内存空间 |
| 通信 | 消息 / 共享内存 | IPC 序列化 | IPC 序列化 |
| 启动开销 | 中等 | 较高 | 较高 |
| 适用 | CPU 密集计算 | 独立程序/沙箱 | HTTP 服务多核 |
| 稳定性 | Worker 崩溃影响同进程 | 子进程崩溃互不影响 | 工作进程崩溃互不影响 |
| 共享状态 | SharedArrayBuffer | 不支持 | 不支持 |
10. 适合与不适合多线程的任务
| 适合多线程 | 不适合多线程 |
|---|---|
| 图像处理/缩略图生成 | 简单的 HTTP 请求路由 |
| 加密/哈希计算 | 数据库 CRUD 操作 |
| 压缩/解压缩 | 文件 I/O(用异步即可) |
| 大规模数学运算 | 简单的 JSON 转换 |
| PDF 生成/渲染 | 短生命周期的微任务 |
| 音视频转码 | 事件驱动的消息转发 |
11. 线程池模式
Worker 创建有一定开销,频繁创建销毁会浪费资源。线程池维护固定数量的 Worker,任务到来时分配空闲 Worker,完成后回收复用。
▶ 示例:自建简易线程池
// pool.js
const { Worker } = require('worker_threads');
class Pool {
constructor(workerFile, size) {
this.workerFile = workerFile;
this.size = size;
this.workers = [];
this.queue = [];
for (let i = 0; i < size; i++) {
const worker = new Worker(workerFile);
worker.busy = false;
worker.on('message', (result) => {
const task = worker.currentTask;
worker.busy = false;
worker.currentTask = null;
task.resolve(result);
this._processQueue();
});
worker.on('error', (err) => {
const task = worker.currentTask;
if (task) {
worker.busy = false;
worker.currentTask = null;
task.reject(err);
this._processQueue();
}
});
this.workers.push(worker);
}
}
run(data) {
return new Promise((resolve, reject) => {
const task = { data, resolve, reject };
const idle = this.workers.find((w) => !w.busy);
if (idle) {
this._assign(idle, task);
} else {
this.queue.push(task);
}
});
}
_assign(worker, task) {
worker.busy = true;
worker.currentTask = task;
worker.postMessage(task.data);
}
_processQueue() {
if (this.queue.length === 0) return;
const idle = this.workers.find((w) => !w.busy);
if (!idle) return;
const task = this.queue.shift();
this._assign(idle, task);
}
destroy() {
for (const worker of this.workers) {
worker.terminate();
}
this.workers = [];
this.queue = [];
}
}
module.exports = Pool;
▶ 示例:使用线程池计算斐波那契
// main.js
const Pool = require('./pool.js');
const pool = new Pool('./fib-worker.js', 4);
async function main() {
const tasks = [40, 41, 42, 43, 44, 45, 46, 47];
const promises = tasks.map((n) => pool.run({ n }));
const results = await Promise.all(promises);
for (let i = 0; i < tasks.length; i++) {
console.log(`fib(${tasks[i]}) = ${results[i]}`);
}
pool.destroy();
}
main();
// fib-worker.js
const { parentPort } = require('worker_threads');
parentPort.on('message', ({ n }) => {
const result = fib(n);
parentPort.postMessage(result);
});
function fib(n) {
if (n <= 1) return n;
let a = 0, b = 1;
for (let i = 2; i <= n; i++) {
[a, b] = [b, a + b];
}
return b;
}
▶ 示例:piscina 线程池库
npm install piscina
const path = require('path');
const Piscina = require('piscina');
const pool = new Piscina({
filename: path.resolve(__dirname, 'task.js'),
maxThreads: 4
});
async function main() {
const results = await Promise.all([
pool.run({ x: 10, y: 20 }),
pool.run({ x: 30, y: 40 }),
pool.run({ x: 50, y: 60 })
]);
console.log(results); // [30, 70, 110]
await pool.destroy();
}
main();
// task.js
module.exports = ({ x, y }) => {
return x + y;
};
12. 综合示例:图像处理线程池
结合前面的知识,构建一个完整的图像处理线程池系统:主线程分发任务、Worker 处理缩略图、结果收集汇总。
// image-pool.js
const { Worker } = require('worker_threads');
const path = require('path');
class ImagePool {
constructor(workerCount) {
this.workers = [];
this.queue = [];
this.results = [];
for (let i = 0; i < workerCount; i++) {
const worker = new Worker(path.join(__dirname, 'image-processor.js'));
worker.busy = false;
worker.on('message', (msg) => {
if (msg.type === 'result') {
this.results.push(msg.data);
}
worker.busy = false;
worker.currentResolve();
this._dispatch();
});
worker.on('error', (err) => {
worker.busy = false;
if (worker.currentReject) {
worker.currentReject(err);
}
this._dispatch();
});
this.workers.push(worker);
}
}
process(fileList) {
this.results = [];
const promises = fileList.map((file) => this._enqueue(file));
return Promise.all(promises).then(() => this.results);
}
_enqueue(file) {
return new Promise((resolve, reject) => {
this.queue.push({ file, resolve, reject });
this._dispatch();
});
}
_dispatch() {
while (this.queue.length > 0) {
const idle = this.workers.find((w) => !w.busy);
if (!idle) break;
const task = this.queue.shift();
idle.busy = true;
idle.currentResolve = task.resolve;
idle.currentReject = task.reject;
idle.postMessage({ file: task.file, sizes: [200, 400, 800] });
}
}
destroy() {
this.workers.forEach((w) => w.terminate());
}
}
module.exports = ImagePool;
// image-processor.js
const { parentPort } = require('worker_threads');
parentPort.on('message', async ({ file, sizes }) => {
const results = [];
for (const size of sizes) {
const thumb = await resize(file, size);
results.push(thumb);
}
parentPort.postMessage({
type: 'result',
data: { file, thumbnails: results }
});
});
async function resize(file, maxSize) {
return new Promise((resolve) => {
const duration = Math.random() * 200 + 50;
setTimeout(() => {
resolve(`${file}_${maxSize}px.jpg`);
}, duration);
});
}
// run.js
const ImagePool = require('./image-pool.js');
async function main() {
const pool = new ImagePool(4);
const files = Array.from({ length: 20 }, (_, i) =>
`photo-${String(i + 1).padStart(3, '0')}.jpg`
);
const start = Date.now();
const results = await pool.process(files);
const elapsed = Date.now() - start;
console.log(`Processed ${files.length} images in ${elapsed} ms`);
console.log(`Generated ${results.reduce((sum, r) => sum + r.thumbnails.length, 0)} thumbnails`);
pool.destroy();
}
main();
node run.js
Processed 20 images in 1234 ms
Generated 60 thumbnails
❓ 常见问题
📖 小节
- 你将学到的核心概念与使用方法
- 故事:Alice 的图像处理加速的核心概念与使用方法
- Worker 类基础的核心概念与使用方法
- parentPort 双向通信的核心概念与使用方法
- workerData 初始化数据的核心概念与使用方法
- MessageChannel 与 MessagePort的核心概念与使用方法
- SharedArrayBuffer 与 Atomics的核心概念与使用方法
- 通信方式对比的核心概念与使用方法
📝 作业
- 完成本课所有代码示例,确保每个示例都能正确运行
- 修改综合示例,添加自己的扩展功能
- 查阅官方文档,找出本课未涉及的1-2个API并编写测试代码
- 思考:在实际项目中,你会如何应用本课学到的知识?
- 尝试将本课知识与前面课程的内容结合,构建一个小项目