Egg.js 多线程间通信实践

2026-07-30 79 浏览 0 评论

在 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 中,进程间通信主要通过以下几种方式:

  1. 内置 Messenger :Egg.js 封装的消息传递机制
  2. 原生 IPC 接口 :直接使用 process.send()process.on('message')
  3. 共享存储 :通过 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 支持通过 onsend 结合实现这一模式。

// 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 评论
点赞
收藏

评论列表 0

暂无评论