18. Cluster 中的资源所有权与重复副作用

本章解决的问题

Cluster 中每个 Worker 都会重新执行入口文件。因此顶层定时器、消息订阅、队列消费者和数据库迁移会执行多次。很多人把它称为“事件重复定义”,但更准确地说:不同进程分别在自己的内存中注册了一次,真正重复的是外部副作用。

为什么会重复

// server.ts 顶层代码
setInterval(cleanExpiredFiles, 60_000);
eventBus.on("task", handleTask);
startQueueConsumer();
runDatabaseMigration();

如果创建四个 Worker,这段入口会执行四次:

Worker 1:自己的 timer、eventBus、consumer
Worker 2:自己的 timer、eventBus、consumer
Worker 3:自己的 timer、eventBus、consumer
Worker 4:自己的 timer、eventBus、consumer

不同进程中的 eventBus 不是同一个 EventEmitter,所以这通常不会直接产生 MaxListenersExceededWarning。但是四个定时器可能同时删除文件,四个消费者可能重复处理不支持竞争消费的任务,四次迁移可能争抢数据库锁。

先定义资源所有者

工作 推荐所有者 原因
Express 路由与请求日志 每个 Worker 每个 Worker 都处理请求
Worker 自己的 DB/Redis 连接 每个 Worker 连接不能跨进程共享
Worker 的优雅关闭监听 每个 Worker 关闭本进程资源
Cluster 整体生命周期 Primary 创建、监控、停止 Worker
数据库迁移 部署步骤 应在服务启动前明确执行一次
全局定时清理 独立调度器/带锁任务 需要单一或幂等执行
可靠后台任务 队列消费者 需要确认、重试和持久化

用分支限定初始化

import * as cluster from "node:cluster";

if (cluster.isPrimary) {
  startClusterManagement();
  startBestEffortCleanupTimer();
} else {
  startHttpServer();
}

放到 Primary 只能保证当前 Cluster Primary 内有一个所有者。它不能解决:

重要任务应使用独立调度器、数据库/Redis 分布式锁,或设计为幂等任务。

幂等比“只执行一次”更可靠

分布式环境中很难永久保证严格只执行一次。更务实的目标是任务可重复执行而结果一致。

创建报告:以 taskId 作为唯一键,重复请求返回同一记录
清理文件:删除前查询状态,文件不存在也视为已完成
更新任务:使用状态条件 UPDATE ... WHERE status = 'pending'
消息消费:记录 messageId,重复消息不重复产生业务副作用

分布式锁也要有:唯一持有者标识、过期时间、续期策略和仅持有者释放的校验。不要用“先 GET 再 SET”这种非原子流程模拟锁。

同一进程内的真正重复监听

热重载、测试重复初始化或循环调用启动函数,可能在同一个 EventEmitter 上重复注册:

function registerProcessEvents(): void {
  process.on("SIGTERM", shutdown);
}

registerProcessEvents();
registerProcessEvents();

解决方式是让初始化入口只调用一次,必要时使用 once(),并在测试中移除监听器。不要通过无限提高 setMaxListeners() 掩盖生命周期错误。

let processEventsRegistered = false;

function registerProcessEvents(): void {
  if (processEventsRegistered) return;
  processEventsRegistered = true;
  process.once("SIGTERM", shutdown);
}

Signal 职责不要混乱

Primary 收到 Signal 时停止扩容并协调 Worker;Worker 收到通知时关闭 HTTP Server 和自己的资源。停机流程必须防重入。

let stopping = false;

async function stopPrimary(): Promise<void> {
  if (stopping) return;
  stopping = true;

  for (const worker of Object.values(cluster.workers ?? {})) {
    worker?.disconnect();
  }
}

process.once("SIGTERM", () => void stopPrimary());

如果 exit 监听器会自动补 Worker,停机时必须检查 stopping,否则 Primary 一边停机一边创建新 Worker。

正常与错误时间线

正常启动:Primary 初始化一次 → fork N 个 Worker → 每个 Worker 初始化 HTTP 资源

错误设计:fork N 个 Worker → 每个入口执行 migration/cron/订阅 → 外部副作用 N 次

正常停机:设置 stopping → disconnect Workers → Worker 清理 → exit → 不再补 Worker

错误停机:Worker exit → exit 监听器无条件 fork → 永远停不下来

跨平台注意事项

常见错误

练习

  1. 创建两个 Worker,在顶层启动每秒日志定时器,观察重复输出并解释它们是否属于同一 EventEmitter。
  2. 把定时器移动到 Primary,再思考两个 PM2 实例时会有几个定时器。
  3. 为“清理 24 小时前临时目录”设计幂等流程和并发保护。
  4. 修改 Worker 重启逻辑:停机状态下不再 fork()

验收清单