Egg.js 多线程间通信实践
在 Node.js 生态中,Egg.js 作为企业级开发框架,继承了 Node.js 的单线程事件循环模型。然而,面对 CPU 密集型任务或高并发场景,单线程模型的局限性逐渐显现。为了充分利用多核 CPU 资源,Egg.js 内置了 Cluster 模式,通过多进程(多线程概念在 Node.js 中通常体现为多进程)架构提升应用的并发处理能力。

本文将深入探讨 Egg.js 多线程(多进程)间的通信机制,详细解析其底层原理与实现方式,并结合具体的代码示例展示实际应用场景。
一、Egg.js 多进程模型概述
1.1 Cluster 模式基础
Egg.js 默认采用 Master-Worker 架构,这一架构基于 Node.js 原生的 cluster 模块实现。在该模式下:
- Master 进程 :负责管理 Worker 进程的生命周期,包括 fork、重启、负载均衡等
- Worker 进程 :实际处理业务逻辑的进程,每个 Worker 都是独立的 Node.js 实例
- Agent 进程 :Egg.js 特有的进程角色,用于处理后台任务、定时任务等公共逻辑
1.2 进程数量配置
Egg.js 中 Worker 进程的数量默认与 CPU 核心数一致,也可以通过配置文件手动指定:
// config/config.default.js
exports.cluster = {
// 工作进程数量,默认等于 CPU 核心数
workers: 4,
};二、多进程间通信的核心机制
2.1 IPC 通信原理
Node.js 的 cluster 模块底层依赖 IPC(Inter-Process Communication,进程间通信)实现 Master 与 Worker 之间的消息传递。IPC 通信基于操作系统提供的管道(Pipe)或 Unix Domain Socket 实现。
在 Egg.js 中,进程间通信主要通过以下几种方式:
- 内置 Messenger :Egg.js 封装的消息传递机制
- 原生 IPC 接口 :直接使用
process.send()和process.on('message') - 共享存储 :通过 Redis、数据库等第三方中间件实现数据共享
2.2 Egg.js Messenger 机制
Egg.js 提供了统一的 Messenger API,屏蔽了底层通信细节,开发者可以通过简单的接口实现进程间消息传递。
Messenger 支持三种消息发送目标:
- app :发送给所有 Worker 进程
- agent :发送给 Agent 进程
- master :发送给 Master 进程
三、实战:多进程间通信实现
3.1 基础消息发送与接收
以下示例展示了如何在 Worker 进程中发送消息给 Agent 进程,并接收返回结果。
// app.js
module.exports = app => {
// 监听来自 agent 的消息
app.messenger.on('agent-to-worker', data => {
console.log(`[Worker ${process.pid}] 收到 agent 消息:`, data);
});
// 向 agent 发送消息
app.messenger.sendToAgent('worker-to-agent', {
from: 'worker',
pid: process.pid,
timestamp: Date.now(),
});
};// agent.js
module.exports = agent => {
// 监听来自 worker 的消息
agent.messenger.on('worker-to-agent', data => {
console.log(`[Agent ${process.pid}] 收到 worker 消息:`, data);
// 向所有 worker 广播消息
agent.messenger.sendToApp('agent-to-worker', {
from: 'agent',
pid: process.pid,
originalData: data,
});
});
};3.2 基于 Action 的请求响应模式
在实际开发中,常常需要类似 HTTP 请求响应的通信模式。Egg.js Messenger 支持通过 on 和 send 结合实现这一模式。
// app/controller/home.js
const Controller = require('egg').Controller;
class HomeController extends Controller {
async getDataFromAgent() {
const { ctx } = this;
// 使用 Promise 封装消息通信
const result = await new Promise((resolve, reject) => {
const requestId = `req_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`;
// 设置超时
const timeout = setTimeout(() => {
ctx.app.messenger.removeAllListeners(`agent-response-${requestId}`);
reject(new Error('Request timeout'));
}, 5000);
// 监听响应
ctx.app.messenger.once(`agent-response-${requestId}`, data => {
clearTimeout(timeout);
resolve(data);
});
// 发送请求
ctx.app.messenger.sendToAgent('worker-request', {
requestId,
action: 'getUserInfo',
payload: { userId: 12345 },
});
});
ctx.body = { success: true, data: result };
}
}
module.exports = HomeController;// agent.js
module.exports = agent => {
agent.messenger.on('worker-request', ({ requestId, action, payload }) => {
// 根据不同 action 处理业务逻辑
let result;
switch (action) {
case 'getUserInfo':
result = {
id: payload.userId,
name: 'Egg.js Developer',
role: 'admin',
lastLogin: new Date().toISOString(),
};
break;
default:
result = { error: 'Unknown action' };
}
// 发送响应
agent.messenger.sendToApp(`agent-response-${requestId}`, result);
});
};3.3 Worker 间直接通信
默认情况下,Worker 进程之间不能直接通信,必须通过 Master 或 Agent 中转。以下示例展示了 Worker A 如何通过 Agent 中转与 Worker B 通信。
// app.js - Worker 发送消息给其他 Worker
module.exports = app => {
// 发送广播消息
app.broadcastToWorkers = function(event, data) {
app.messenger.sendToAgent('broadcast-request', {
event,
data,
excludePid: process.pid,
});
};
// 监听广播消息
app.messenger.on('broadcast-message', ({ event, data, fromPid }) => {
console.log(`[Worker ${process.pid}] 收到来自 ${fromPid} 的广播 [${event}]:`, data);
});
};// agent.js - 中转广播
module.exports = agent => {
agent.messenger.on('broadcast-request', ({ event, data, excludePid }) => {
// sendToApp 会发送给所有 worker
agent.messenger.sendToApp('broadcast-message', {
event,
data,
fromPid: excludePid,
});
});
};四、共享数据方案
4.1 使用 Redis 实现共享状态
对于需要在多进程间共享的状态数据,Redis 是常用的解决方案。以下示例展示了基于 Redis 的计数器实现。
// app/service/counter.js
const Service = require('egg').Service;
class CounterService extends Service {
async increment(key, delta = 1) {
const { redis } = this.app;
const newValue = await redis.incrby(key, delta);
return newValue;
}
async getValue(key) {
const { redis } = this.app;
const value = await redis.get(key);
return parseInt(value, 10) || 0;
}
async reset(key) {
const { redis } = this.app;
await redis.set(key, 0);
return 0;
}
}
module.exports = CounterService;4.2 基于消息的缓存同步
当某个 Worker 更新了本地缓存后,需要通知其他 Worker 同步失效,这一场景可以通过消息机制实现。
// app.js
module.exports = app => {
// 本地缓存对象
const localCache = new Map();
// 缓存失效通知
app.messenger.on('cache-invalidate', ({ key }) => {
if (localCache.has(key)) {
localCache.delete(key);
console.log(`[Worker ${process.pid}] 缓存已失效: ${key}`);
}
});
// 扩展 app 方法
app.invalidateCache = function(key) {
// 先清除本地
localCache.delete(key);
// 通知其他进程
app.messenger.sendToAgent('cache-invalidate-broadcast', { key });
};
app.getFromCache = function(key) {
return localCache.get(key);
};
app.setToCache = function(key, value) {
localCache.set(key, value);
};
};// agent.js
module.exports = agent => {
// 中转缓存失效广播
agent.messenger.on('cache-invalidate-broadcast', ({ key }) => {
agent.messenger.sendToApp('cache-invalidate', { key });
});
};五、性能考量与实现细节
5.1 消息序列化
进程间通信传递的数据需要经过序列化和反序列化。Node.js IPC 默认使用 JSON 序列化,因此传递的数据必须是可序列化的对象,不能包含函数、循环引用等。
// 正确:可序列化的数据
app.messenger.sendToAgent('valid-message', {
name: 'test',
count: 100,
list: [1, 2, 3],
nested: { a: 1, b: 2 },
});
// 错误:不可序列化的数据(包含函数)
// app.messenger.sendToAgent('invalid-message', {
// name: 'test',
// callback: function() { ... }
// });5.2 消息大小限制
IPC 通信存在消息大小限制,通常单条消息不宜过大。对于大数据量的传递,建议采用分片传输或通过共享文件、Redis 等方式中转。
// 大数据分片传输示例
async function sendLargeData(app, data) {
const CHUNK_SIZE = 64 * 1024; // 64KB per chunk
const totalSize = Buffer.byteLength(JSON.stringify(data));
const chunks = Math.ceil(totalSize / CHUNK_SIZE);
const sessionId = `large_${Date.now()}`;
// 发送开始通知
app.messenger.sendToAgent('large-data-start', { sessionId, chunks });
// 分片发送
const jsonStr = JSON.stringify(data);
for (let i = 0; i < chunks; i++) {
const chunk = jsonStr.slice(i * CHUNK_SIZE, (i + 1) * CHUNK_SIZE);
app.messenger.sendToAgent('large-data-chunk', {
sessionId,
index: i,
chunk,
});
}
}5.3 消息丢失与重发机制
在进程重启等异常场景下,消息可能丢失。对于重要消息,可以实现确认重发机制。
// 带确认的消息发送
function sendWithAck(app, event, data, timeout = 3000, maxRetries = 3) {
return new Promise((resolve, reject) => {
let retries = 0;
let timer;
const messageId = `msg_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`;
const ackHandler = ackData => {
if (ackData.messageId === messageId) {
clearTimeout(timer);
app.messenger.removeListener(`${event}-ack`, ackHandler);
resolve(ackData.result);
}
};
const send = () => {
if (retries >= maxRetries) {
app.messenger.removeListener(`${event}-ack`, ackHandler);
reject(new Error(`Message ${event} failed after ${maxRetries} retries`));
return;
}
retries++;
app.messenger.sendToAgent(event, { messageId, data });
timer = setTimeout(() => {
console.warn(`[Worker ${process.pid}] 消息超时,第 ${retries} 次重发: ${event}`);
send();
}, timeout);
};
app.messenger.on(`${event}-ack`, ackHandler);
send();
});
}六、典型应用场景
6.1 分布式任务调度
在多 Worker 场景下,定时任务如果每个 Worker 都执行会导致重复执行。通过 Agent 进程统一调度,可以避免这一问题。
// agent.js
module.exports = agent => {
// Agent 中执行定时任务
const schedule = require('node-schedule');
// 每小时执行一次的数据同步任务
schedule.scheduleJob('0 0 * * * *', () => {
console.log('[Agent] 开始执行数据同步任务');
// 模拟数据处理
const syncData = {
taskId: `sync_${Date.now()}`,
syncTime: new Date().toISOString(),
records: 1500,
};
// 将结果分发给所有 Worker
agent.messenger.sendToApp('data-synced', syncData);
});
};// app.js
module.exports = app => {
// 接收 Agent 同步的数据
app.messenger.on('data-synced', data => {
console.log(`[Worker ${process.pid}] 数据已同步:`, data.taskId);
// 更新本地缓存或状态
app.locals.syncData = data;
});
};6.2 进程健康监控
通过 Master 与 Worker 的通信,可以实现进程健康状态的实时监控。
// app.js
module.exports = app => {
// 定期向 Agent 上报健康状态
setInterval(() => {
const usage = process.memoryUsage();
app.messenger.sendToAgent('worker-health-report', {
pid: process.pid,
memory: {
rss: Math.round(usage.rss / 1024 / 1024), // MB
heapUsed: Math.round(usage.heapUsed / 1024 / 1024),
heapTotal: Math.round(usage.heapTotal / 1024 / 1024),
},
uptime: process.uptime(),
activeRequests: app.server ? app.server._connections : 0,
});
}, 30000); // 每 30 秒上报一次
};// agent.js
module.exports = agent => {
// 存储各 Worker 的健康状态
const workerHealth = new Map();
agent.messenger.on('worker-health-report', report => {
workerHealth.set(report.pid, {
...report,
reportTime: Date.now(),
});
});
// 提供健康状态查询接口
agent.messenger.on('get-health-status', ({ requestId }) => {
const status = Array.from(workerHealth.entries()).map(([pid, data]) => ({
pid,
...data,
}));
agent.messenger.sendToApp(`health-status-${requestId}`, status);
});
};七、总结
Egg.js 的多进程通信机制建立在 Node.js Cluster 模块的 IPC 通信基础之上,通过 Messenger 封装提供了简洁易用的 API。在实际开发中,开发者需要根据业务场景选择合适的通信方式:
- 简单的指令传递可直接使用 Messenger API
- 请求响应式通信需要自行封装 Promise 和超时机制
- 大数据共享建议采用 Redis 等外部存储
- 定时任务、公共逻辑应放置在 Agent 进程中统一处理
理解多进程通信的底层原理,合理设计通信架构,能够有效提升 Egg.js 应用的稳定性和性能表现。
发布评论
评论列表 0



