文章

进程与线程深度解析

进程与线程深度解析

一句话概括

Node.js 虽然主线程是单线程的,但通过 child_processclusterworker_threads 分别提供了多进程、多实例和多线程的并发能力,本文将深入剖析这三者的原理、适用场景以及如何在 Node.js 中有效利用多核 CPU 和并行计算。

背景与意义

「Node.js 是单线程的」这句话常常被误解为「Node.js 只能在一个 CPU 核心上运行」。实际上,Node.js 的单线程指的是 JavaScript 代码执行在主线程上,但这并不意味着它无法利用多核 CPU。

想象一个电商网站的 Node.js 服务器:作为主线程需要处理请求路由、业务逻辑、数据库查询调度等工作。如果只有一个进程在跑,即使服务器有 32 核 CPU,也只能用到其中 1 核,这显然是一种极大的浪费。

Node.js 提供了三层并行方案:

  • child_process:适合需要独立运行的子任务
  • cluster:适合充分利用多核 CPU 的 HTTP 服务
  • worker_threads:适合 CPU 密集型计算

理解这三者的区别和联系,是 Node.js 后端架构设计的必修课。

概念与定义

进程(Process):操作系统资源分配的基本单位,每个进程有独立的地址空间、内存、文件描述符。进程间通信(IPC)需要通过管道、消息队列、共享内存等机制。

线程(Thread):CPU 调度的基本单位,同一进程内的线程共享进程的地址空间和资源,线程间通信更高效。

child_process:Node.js 内置模块,用于创建子进程并与其通信。可以执行系统命令、运行其他 Node.js 脚本或其他语言的可执行文件。

cluster:基于 child_process 封装的模块,专为 HTTP 服务设计,自动将请求负载均衡到多个工作进程。

worker_threads:Node.js v10.5.0 引入,提供真正的线程级并行,同一进程内的多个线程共享内存(通过 SharedArrayBuffer)。

核心知识点拆解

1. child_process:创建子进程的四种方式

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
const { spawn, exec, execFile, fork } = require('child_process');
const path = require('path');

/**
 * 1. spawn - 流式输出,适合大量数据
 * 返回 ChildProcess 对象,通过 stream 处理输出
 * 最底层的 API,推荐优先使用
 */
function spawnExample() {
  // 执行系统命令,列出当前目录
  const child = spawn('ls', ['-la', __dirname], {
    stdio: ['pipe', 'pipe', 'pipe'], // 标准输入/输出/错误
    cwd: __dirname,  // 工作目录
    env: { NODE_ENV: 'development' }, // 环境变量
    maxBuffer: 1024 * 1024, // 最大输出缓冲区
  });
  
  child.stdout.on('data', (data) => {
    console.log(`标准输出:\n${data}`);
  });
  
  child.stderr.on('data', (data) => {
    console.error(`标准错误:\n${data}`);
  });
  
  child.on('close', (code, signal) => {
    console.log(`子进程退出,code=${code}, signal=${signal}`);
  });
  
  child.on('error', (err) => {
    console.error('无法启动子进程:', err);
  });
}

/**
 * 2. exec - 缓冲输出,适合少量数据
 * 将输出全部缓冲到内存后通过 callback 一次性返回
 */
function execExample() {
  exec('ls -la', { maxBuffer: 1024 * 1024 }, (error, stdout, stderr) => {
    if (error) {
      console.error(`执行出错: ${error.message}`);
      return;
    }
    if (stderr) {
      console.error(`stderr: ${stderr}`);
    }
    console.log(`stdout:\n${stdout}`);
  });
}

/**
 * 3. execFile - 直接执行文件,不通过 shell
 * 比 exec 更安全(避免 shell 注入),性能更好
 */
function execFileExample() {
  execFile('node', ['--version'], (error, stdout, stderr) => {
    if (error) throw error;
    console.log(`Node.js 版本: ${stdout.trim()}`);
  });
}

/**
 * 4. fork - 专为 Node.js 子进程设计
 * 自带 IPC 通道,方便进程间通信
 */
function forkExample() {
  const child = fork(path.join(__dirname, 'worker.js'), ['--arg1', 'value1'], {
    silent: true, // 不继承父进程的 stdio
    execArgv: ['--max-old-space-size=256'], // V8 参数
  });
  
  // 发送消息到子进程
  child.send({ type: 'task', payload: { id: 1, data: 'hello' } });
  
  // 接收子进程的消息
  child.on('message', (msg) => {
    console.log('收到子进程消息:', msg);
  });
  
  child.on('exit', (code) => {
    console.log(`子进程退出,code=${code}`);
  });
}

子进程通信(IPC)的高级用法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
// parent.js
const { fork } = require('child_process');

const workers = [];
const TASK_QUEUE = [];
const RESULT_MAP = new Map();

// 创建子进程池
function createWorkerPool(size) {
  for (let i = 0; i < size; i++) {
    const worker = fork('./worker.js');
    worker.id = i;
    worker.busy = false;
    
    worker.on('message', (msg) => {
      worker.busy = false;
      
      const resolve = RESULT_MAP.get(msg.taskId);
      if (resolve) {
        resolve(msg.result);
        RESULT_MAP.delete(msg.taskId);
      }
      
      assignNextTask(worker);
    });
    
    workers.push(worker);
  }
}

function assignNextTask(worker) {
  if (TASK_QUEUE.length > 0 && !worker.busy) {
    const task = TASK_QUEUE.shift();
    worker.busy = true;
    worker.send(task);
  }
}

function submitTask(task) {
  return new Promise((resolve) => {
    const taskId = `${Date.now()}-${Math.random()}`;
    RESULT_MAP.set(taskId, resolve);
    
    const taskWithId = { ...task, taskId };
    
    // 找空闲 worker
    const idleWorker = workers.find(w => !w.busy);
    if (idleWorker) {
      idleWorker.busy = true;
      idleWorker.send(taskWithId);
    } else {
      TASK_QUEUE.push(taskWithId);
    }
  });
}

// worker.js
process.on('message', (msg) => {
  // 执行任务(假设是 CPU 密集型任务)
  const result = performHeavyComputation(msg.payload);
  
  // 发送结果回父进程
  process.send({ taskId: msg.taskId, result });
});

function performHeavyComputation(data) {
  // 模拟耗时计算
  let result = 0;
  for (let i = 0; i < 5e7; i++) {
    result += Math.sqrt(i) * Math.sin(i);
  }
  return result;
}

2. cluster:多核 HTTP 服务

cluster 是 Node.js 中最常用的多核利用方案。它基于 child_process.fork(),但封装了端口共享和负载均衡。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
const cluster = require('cluster');
const http = require('http');
const os = require('os');

/**
 * 基本 cluster 服务器
 * 主进程负责调度,工作进程处理请求
 */
if (cluster.isMaster) {
  const numCPUs = os.cpus().length;
  console.log(`主进程 ${process.pid} 启动,CPU 核心数: ${numCPUs}`);
  
  // 创建工作进程
  for (let i = 0; i < numCPUs; i++) {
    cluster.fork();
  }
  
  // 监听工作进程退出,自动重启
  cluster.on('exit', (worker, code, signal) => {
    console.log(`工作进程 ${worker.process.pid} 已退出`);
    // 自动重启
    console.log('启动新的工作进程...');
    cluster.fork();
  });
  
  // 优雅关闭
  process.on('SIGTERM', () => {
    for (const id in cluster.workers) {
      cluster.workers[id].kill('SIGTERM');
    }
    process.exit(0);
  });
  
} else {
  // 工作进程:处理 HTTP 请求
  const server = http.createServer((req, res) => {
    res.writeHead(200);
    res.end(`处理请求的工作进程: ${process.pid}\n`);
  });
  
  server.listen(8000);
  console.log(`工作进程 ${process.pid} 已启动`);
}

cluster 的高级配置:调度策略与粘性会话

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
const cluster = require('cluster');
const http = require('http');

if (cluster.isMaster) {
  // 调度策略:默认是 round-robin(轮询)
  // cluster.SCHED_RR - 轮询(Windows 默认)
  // cluster.SCHED_NONE - 操作系统调度(macOS/Unix 默认)
  cluster.schedulingPolicy = cluster.SCHED_RR;
  
  const worker = cluster.fork();
  
  // 主进程可以向工作进程发送消息
  worker.send({ type: 'config', data: { port: 3000 } });
  
  // 也可以监听工作进程的消息
  worker.on('message', (msg) => {
    if (msg.type === 'metrics') {
      console.log(`工作进程 ${worker.id} 指标:`, msg.data);
    }
  });
  
} else {
  // 工作进程健康检查
  setInterval(() => {
    const memoryUsage = process.memoryUsage();
    process.send({
      type: 'metrics',
      data: {
        pid: process.pid,
        memory: memoryUsage,
        uptime: process.uptime(),
      }
    });
  }, 60000);
}

/**
 * sticky session(粘性会话)实现
 * 对于需要保持 session 的应用,需要将同一用户的请求发送到同一个 worker
 */
const sticky = require('sticky-session'); // 或使用 nginx ip_hash

// 或者自定义实现:
function createStickyServer(port) {
  if (cluster.isMaster) {
    const workers = [];
    
    for (let i = 0; i < os.cpus().length; i++) {
      workers.push(cluster.fork());
    }
    
    // 创建一个「转发服务器」根据 IP 分配 worker
    const proxy = require('net').createServer({ pauseOnConnect: true }, (socket) => {
      // 基于客户端 IP 哈希选择 worker
      const ipHash = socket.remoteAddress
        .split('')
        .reduce((acc, char) => acc + char.charCodeAt(0), 0);
      const workerId = ipHash % workers.length;
      
      // 将连接发送给选中的 worker
      workers[workerId].send('connection', socket);
    });
    
    proxy.listen(port);
  } else {
    // 工作进程处理 HTTP
    // ...正常 HTTP 服务代码
  }
}

3. worker_threads:真正的多线程并行

worker_threads 是 Node.js 真正的线程级并行的解决方案。与 child_process 不同,线程共享同一进程的内存空间,通信成本更低。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
const { Worker, isMainThread, parentPort, workerData } = require('worker_threads');
const path = require('path');

/**
 * 主线程:创建工作线程池
 */

if (isMainThread) {
  class ThreadPool {
    constructor(workerScript, poolSize = 4) {
      this.workers = [];
      this.available = [];
      this.taskQueue = [];
      
      for (let i = 0; i < poolSize; i++) {
        this.createWorker(workerScript);
      }
    }
    
    createWorker(workerScript) {
      const worker = new Worker(path.resolve(__dirname, workerScript), {
        workerData: { workerId: this.workers.length },
        // 可以传递初始数据
      });
      
      worker.on('message', (result) => {
        // 找到对应的 resolve 并调用
        if (worker.currentResolve) {
          worker.currentResolve(result);
          worker.currentResolve = null;
        }
        
        // 处理下一个任务
        if (this.taskQueue.length > 0) {
          const { data, resolve } = this.taskQueue.shift();
          worker.currentResolve = resolve;
          worker.postMessage(data);
        } else {
          this.available.push(worker);
        }
      });
      
      worker.on('error', (err) => {
        console.error('Worker error:', err);
      });
      
      worker.on('exit', (code) => {
        if (code !== 0) {
          console.error(`Worker exited with code ${code}`);
          // 重启
          this.createWorker(workerScript);
        }
      });
      
      this.workers.push(worker);
      this.available.push(worker);
    }
    
    execute(data) {
      return new Promise((resolve) => {
        if (this.available.length > 0) {
          const worker = this.available.pop();
          worker.currentResolve = resolve;
          worker.postMessage(data);
        } else {
          this.taskQueue.push({ data, resolve });
        }
      });
    }
    
    async terminate() {
      await Promise.all(this.workers.map(w => w.terminate()));
    }
  }
  
  // 使用线程池
  async function main() {
    const pool = new ThreadPool('./heavy-worker.js', 4);
    
    const tasks = Array.from({ length: 10 }, (_, i) => ({
      id: i,
      data: Array.from({ length: 1000 }, () => Math.random() * 1000)
    }));
    
    const results = await Promise.all(
      tasks.map(task => pool.execute(task))
    );
    
    console.log('所有任务完成:', results.length);
    await pool.terminate();
  }
  
  main().catch(console.error);
}

/**
 * 工作线程:heavy-worker.js
 * 在单独的线程中执行 CPU 密集型计算
 */
// if (!isMainThread) {
//   const { parentPort, workerData } = require('worker_threads');
  
//   parentPort.on('message', (task) => {
//     const { id, data } = task;
    
//     // 执行 CPU 密集型任务
//     const result = {
//       taskId: id,
//       sum: data.reduce((a, b) => a + b, 0),
//       average: data.reduce((a, b) => a + b, 0) / data.length,
//       sorted: data.sort((a, b) => a - b),
//       median: calculateMedian(data),
//       workerId: workerData.workerId,
//     };
    
//     // 将结果发送回主线程
//     parentPort.postMessage(result);
//   });
  
//   function calculateMedian(sorted) {
//     const mid = Math.floor(sorted.length / 2);
//     return sorted.length % 2 === 0
//       ? (sorted[mid - 1] + sorted[mid]) / 2
//       : sorted[mid];
//   }
// }

SharedArrayBuffer:线程间共享内存

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
const { Worker, isMainThread, parentPort, workerData } = require('worker_threads');

if (isMainThread) {
  // 创建共享内存缓冲区
  const sharedBuffer = new SharedArrayBuffer(4 * 1024); // 4KB 共享内存
  const sharedArray = new Int32Array(sharedBuffer);
  sharedArray[0] = 0; // 初始值
  
  // 创建 4 个工作线程,共享同一内存
  const workers = [];
  for (let i = 0; i < 4; i++) {
    const worker = new Worker(__filename, {
      workerData: { sharedBuffer, workerId: i }
    });
    workers.push(worker);
  }
  
  // 等待所有线程完成
  Promise.all(workers.map(w => new Promise(resolve => {
    w.on('exit', resolve);
  }))).then(() => {
    console.log('最终计数器值:', sharedArray[0]);
  });
  
} else {
  // 工作线程代码
  const { sharedBuffer, workerId } = workerData;
  const sharedArray = new Int32Array(sharedBuffer);
  
  // 使用 Atomics 进行线程安全操作
  for (let i = 0; i < 1000000; i++) {
    Atomics.add(sharedArray, 0, 1);
  }
  
  console.log(`Worker ${workerId} 完成`);
}

实战案例

实战:图片处理服务

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
const { Worker } = require('worker_threads');
const http = require('http');
const { Transform } = require('stream');

/**
 * 使用 worker_threads 处理图片缩略图生成
 * 将 CPU 密集型的图片处理任务移出主线程
 */

// image-worker.js
// const { parentPort, workerData } = require('worker_threads');
// const sharp = require('sharp'); // 假设已安装 sharp

// parentPort.on('message', async ({ id, imageBuffer, options }) => {
//   try {
//     const result = await sharp(imageBuffer)
//       .resize(options.width, options.height)
//       .jpeg({ quality: 80 })
//       .toBuffer();
    
//     parentPort.postMessage({ id, result, size: result.length });
//   } catch (error) {
//     parentPort.postMessage({ id, error: error.message });
//   }
// });

// 主服务器
const workerPool = [];
const MAX_WORKERS = 4;

for (let i = 0; i < MAX_WORKERS; i++) {
  const worker = new Worker('./image-worker.js');
  worker.busy = false;
  worker.pending = [];
  workerPool.push(worker);
  
  worker.on('message', (msg) => {
    worker.busy = false;
    if (worker.currentCallback) {
      worker.currentCallback(msg);
      worker.currentCallback = null;
    }
    // 处理队列中的下一个任务
    if (worker.pending.length > 0) {
      const next = worker.pending.shift();
      worker.busy = true;
      worker.currentCallback = next.callback;
      worker.postMessage(next.data);
    }
  });
}

function getAvailableWorker() {
  return workerPool.find(w => !w.busy) || workerPool[0];
}

const server = http.createServer((req, res) => {
  if (req.url.startsWith('/thumbnail')) {
    // 假设从请求中获取图片数据
    const imageBuffer = Buffer.alloc(1024 * 1024); // 示例
    
    const worker = getAvailableWorker();
    const taskData = {
      id: Date.now(),
      imageBuffer,
      options: { width: 200, height: 200 }
    };
    
    if (!worker.busy) {
      worker.busy = true;
      worker.currentCallback = (result) => {
        if (result.error) {
          res.writeHead(500);
          res.end(result.error);
        } else {
          res.writeHead(200, { 'Content-Type': 'image/jpeg' });
          res.end(result.result);
        }
      };
      worker.postMessage(taskData);
    } else {
      worker.pending.push({
        data: taskData,
        callback: (result) => {
          if (result.error) {
            res.writeHead(500);
            res.end(result.error);
          } else {
            res.writeHead(200, { 'Content-Type': 'image/jpeg' });
            res.end(result.result);
          }
        }
      });
    }
  }
});

server.listen(3000);

底层原理

child_process 的内部实现

child_process.spawn() 的底层依赖 libuv 的 uv_spawn()。在 Unix 系统上,其流程是:

1
2
3
4
5
6
7
8
9
10
11
1. 创建 IPC 管道(socket pair)
2. fork() 子进程
3. 在子进程中:
   a. 关闭不需要的文件描述符
   b. 设置 stdio 重定向(pipe/ignore/inherit)
   c. 设置工作目录、环境变量
   d. 调用 execve() 执行目标程序
4. 父进程中:
   a. 创建 ChildProcess 对象
   b. 建立 IPC 通信
   c. 返回 ChildProcess 引用

cluster 的端口共享机制

cluster 的实现依赖于一个关键操作:文件描述符传递

1
2
3
4
5
6
7
8
9
10
11
12
13
主进程:
  1. 创建 TCP 服务器(listen 端口)
  2. 不直接 accept 连接
  3. 将服务器文件描述符通过 IPC 发送给工作进程

工作进程:
  1. 收到文件描述符
  2. 创建 HTTP 服务器绑定到该描述符
  3. 多个工作进程可以同时 accept 同一个端口

负载均衡(cluster.SCHED_RR):
  主进程 accept 连接 → 按 RR 算法分配给工作进程
  工作进程直接处理请求

child_process vs worker_threads 的性能对比

维度child_processworker_threads
内存独立进程空间,内存开销大共享进程空间,内存开销小
启动速度慢(需要创建新进程)快(仅创建线程)
通信开销较高(IPC 序列化/反序列化)较低(共享内存 + 消息传递)
隔离性完全隔离,一个崩溃不影响其他不隔离,一个线程崩溃可能导致整个进程崩溃
适合场景系统命令执行、隔离任务CPU 密集型计算

高频面试题解析

面试题1:port 冲突——为什么 cluster 下多个 worker 可以监听同一端口?

:cluster 模式下,实际的主进程(master)监听端口,然后通过文件描述符传递将 socket 发送给工作进程。多个工作进程共享同一个文件描述符,所以它们可以同时 accept 同一端口上的连接。操作系统内核级别的 SO_REUSEPORT(Linux 3.9+)或 master 进程层面的分发实现了这一点。

面试题2:process.exit() 和 worker.terminate() 的区别?

process.exit() 终止整个进程,所有工作线程会被强制终止。worker.terminate() 只终止特定工作线程,主进程和其他线程正常运行。使用 terminate() 时,工作线程可能正在处理任务,所以需要配套任务恢复或重试机制。

面试题3:child_process 的 spawn 和 exec 各有什么优缺点?

 spawnexec
I/O 方式流式(可以处理大量数据)缓冲(限制 maxBuffer)
Shell默认不使用使用 shell
安全性无 shell 注入风险有 shell 注入风险
回调事件监听callback 一次性返回
适用场景长时间运行、大输出量简单命令、小输出量

面试题4:worker_threads 能提高 I/O 密集型应用的性能吗?

:通常不能。Node.js 的事件循环已经可以高效处理 I/O 密集型操作。worker_threads 的设计目标是 CPU 密集型任务。对于 I/O 密集型,增加线程反而会增加上下文切换的开销。Node.js 本身的单线程事件循环 + libuv 线程池已经是 I/O 密集型场景的绝佳方案。

面试题5:如何优雅地关闭 cluster 应用?

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
if (cluster.isMaster) {
  process.on('SIGTERM', () => {
    for (const id in cluster.workers) {
      // 发送关闭信号
      cluster.workers[id].send('shutdown');
      // 设置超时,强行走
      setTimeout(() => {
        cluster.workers[id].kill('SIGKILL');
      }, 30000);
    }
  });
} else {
  process.on('message', (msg) => {
    if (msg === 'shutdown') {
      server.close(() => {
        // 不再接受新请求,处理完当前请求后退出
        process.exit(0);
      });
      
      // 设置超时,强制退出
      setTimeout(() => process.exit(1), 20000);
    }
  });
}

总结与扩展

模块本质通信方式资源开销隔离性输入输出类型
child_process新进程IPC完整隔离I/O 为主
cluster同进程的多实例IPC + socket 共享较高进程级隔离HTTP 服务
worker_threads同一进程内的线程消息传递 + 共享内存不隔离CPU 密集型

扩展方向

  • PM2 进程管理器:PM2 是对 cluster 模块的高层封装,提供了进程守护、日志管理、零停机重启等功能
  • Docker + Node.js:在容器化环境中,通常每个容器运行一个进程,通过 Kubernetes 等进行实例扩缩
  • NAPI 与 C++ 插件:对于极端性能要求,可以编写 C++ Addon,绕过线程/进程通信开销
  • 单进程性能调优:在不能使用多进程的场景下,优化事件循环和 libuv 线程池的利用率
本文由作者按照 CC BY 4.0 进行授权

© 独行的风. 保留部分权利。

本站采用 Jekyll 主题 Chirpy

本站总访问量 本站访客数 本文阅读量