进程与线程深度解析
一句话概括
Node.js 虽然主线程是单线程的,但通过 child_process、cluster 和 worker_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_process | worker_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 各有什么优缺点?
| spawn | exec | |
|---|---|---|
| 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 线程池的利用率