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 内有一个所有者。它不能解决:
- 同时部署多个应用实例;
- Primary 重启后定时状态丢失;
- 清理执行一半时崩溃;
- Kubernetes 中多个 Pod 同时运行。
重要任务应使用独立调度器、数据库/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 → 永远停不下来
跨平台注意事项
- 依赖
flock等文件锁的方案通常不能直接移植到 Windows。 - Windows 和 Linux 的 Signal 触发方式不同,不能把 Ctrl+C 测试当成完整的生产停机测试。
- 网络文件系统的锁与原子操作语义可能不同于本地磁盘。
- 跨主机协调应依赖共享数据库、Redis 或队列,不依赖某台机器的 Primary。
常见错误
- 把“每个进程各有一个监听器”误判为同一 EventEmitter 的监听泄漏。
- 只在 Primary 启动定时器,就宣称整个集群严格只运行一次。
- 数据库迁移跟随每个 Worker 启动。
- 用无过期时间的锁,持有进程崩溃后永远无法恢复。
- 定时任务不幂等,重试就产生重复数据。
- 停机期间
exit回调仍不断补 Worker。
练习
- 创建两个 Worker,在顶层启动每秒日志定时器,观察重复输出并解释它们是否属于同一 EventEmitter。
- 把定时器移动到 Primary,再思考两个 PM2 实例时会有几个定时器。
- 为“清理 24 小时前临时目录”设计幂等流程和并发保护。
- 修改 Worker 重启逻辑:停机状态下不再
fork()。
验收清单
- [ ] 能解释入口文件为何会在每个 Worker 中执行。
- [ ] 能区分跨进程重复副作用与单进程重复监听。
- [ ] 每项后台工作都有明确所有者。
- [ ] 关键任务是幂等的,并有可靠协调机制。
- [ ] 数据库迁移不跟随每个 Worker 启动。
- [ ] 停机状态下不会自动补充 Worker。