Node.js: Worker 线程

最后更新:2026-08-26

1. 你将学到



2. 故事:Alice 的图像处理加速

Alice 在一家 SaaS 公司负责图像处理服务。用户上传图片后,系统需要为每张图片生成 3 种尺寸的缩略图。最初的单线程方案处理 100 张图片需要约 30 秒,高峰期请求排队严重。

她决定用 Worker Threads 将任务分配到 4 个线程并行处理:主线程读取文件列表,通过 workerData 传入任务参数,Worker 完成缩略图生成后通过 parentPort 返回结果。最终 100 张图片的处理时间降到约 8 秒,吞吐量提升近 4 倍。



3. Worker 类基础

worker_threads 模块是 Node.js 提供的真正的多线程能力。每个 Worker 在独立线程中执行脚本,拥有自己的 V8 引擎实例和事件循环。

JAVASCRIPT
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 脚本内部:

JAVASCRIPT
// 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 创建与通信

JAVASCRIPT
// 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 事件接收数据。

▶ 示例:双向消息交互

JAVASCRIPT
// 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 });
  }
});
▶ 试一试
JAVASCRIPT
// 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 对象零拷贝传输

JAVASCRIPT
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');
▶ 试一试
JAVASCRIPT
// 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

JAVASCRIPT
// 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);
  });
}
▶ 试一试
JAVASCRIPT
// 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,实现线程间直接通信,无需经过主线程中转。

100%
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 直接通信

JAVASCRIPT
// 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]);
▶ 试一试
JAVASCRIPT
// 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 提供的原子操作,可安全地读写共享数据,避免竞态条件。

▶ 示例:共享计数器

JAVASCRIPT
// 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}`);
  });
}
▶ 试一试
JAVASCRIPT
// 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,完成后回收复用。

▶ 示例:自建简易线程池

JAVASCRIPT 📖 仅展示
// 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;
逻辑代码 61 行(超过 40 行限制,仅展示)

▶ 示例:使用线程池计算斐波那契

JAVASCRIPT
// 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();
▶ 试一试
JAVASCRIPT
// 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 线程池库

BASH
npm install piscina
JAVASCRIPT
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();
JAVASCRIPT
// task.js
module.exports = ({ x, y }) => {
  return x + y;
};


12. 综合示例:图像处理线程池

结合前面的知识,构建一个完整的图像处理线程池系统:主线程分发任务、Worker 处理缩略图、结果收集汇总。

JAVASCRIPT
// 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;
JAVASCRIPT
// 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);
  });
}
JAVASCRIPT
// 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();
BASH
node run.js
TEXT 📖 仅展示
Processed 20 images in 1234 ms
Generated 60 thumbnails

❓ 常见问题

Q Worker Threads 和 child_process 有什么区别?
A Worker Threads 在同一进程内创建线程,可共享内存(SharedArrayBuffer),通信开销低;child_process 创建独立进程,内存隔离,通过 IPC 序列化通信,开销较高但更安全。
Q Worker 能访问主线程的变量吗?
A 不能。Worker 有独立的执行上下文,不能直接访问主线程变量。需要通过 postMessage 消息传递、workerData 初始化传值,或 SharedArrayBuffer 共享内存。
Q 什么时候应该用 Worker Threads?
A CPU 密集型任务如图像处理、加密计算、压缩解压、大规模数学运算、音视频转码等。I/O 密集型任务不需要 Worker,Node.js 的异步 I/O 已经足够高效。
Q SharedArrayBuffer 安全吗?
A SharedArrayBuffer 本身不提供同步机制,多线程同时读写可能导致竞态条件。必须配合 Atomics 操作(add、load、store、compareExchange 等)保证原子性,或使用锁模式协调访问。
Q Worker 的创建开销大吗?
A 创建 Worker 需要初始化 V8 引擎实例和事件循环,开销约为 30-50 ms,比 child_process 小但仍然不可忽略。建议使用线程池模式复用 Worker,避免频繁创建和销毁。
Q Worker 内可以使用 fs、http 等 Node.js 模块吗?
A 可以。Worker 线程拥有完整的 Node.js 运行时,支持绝大多数内置模块。但某些模块如 cluster 只能在主进程使用。
Q Transferable 和结构化克隆有什么区别?
A 结构化克隆会复制数据,适合小数据;Transferable 将所有权转移给目标线程,零拷贝但原线程失去访问权,适合大 ArrayBuffer 等场景。

📖 小节


📝 作业

  1. 完成本课所有代码示例,确保每个示例都能正确运行
  2. 修改综合示例,添加自己的扩展功能
  3. 查阅官方文档,找出本课未涉及的1-2个API并编写测试代码
  4. 思考:在实际项目中,你会如何应用本课学到的知识?
  5. 尝试将本课知识与前面课程的内容结合,构建一个小项目
Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏