Node.js后台进程优雅关闭如何确保Nats事件全部处理完成
Node.js 服务搭配NATS的优雅关闭实现方案
核心逻辑本质是:收到关闭信号后先停止接入新任务,等待所有已接收任务全部执行完成后,再释放资源退出进程。你当前的实现只覆盖了HTTP请求的等待逻辑,没有将NATS消息处理、后台任务纳入关闭等待的统计范围,且nc.drain()本身仅负责将NATS服务端缓存的未投递消息全部拉取到客户端、取消订阅停止接收新消息,不会感知业务逻辑的执行状态,因此会出现任务被强制中断的问题。
具体实现步骤
1. 实现轻量全局任务追踪器
不需要引入额外依赖,通过计数器+Promise实现待完成任务的状态追踪即可:所有异步任务(HTTP请求处理、NATS消息处理、自定义后台任务)启动时增加计数,任务执行完成(无论成功/失败)时减少计数,当关闭流程触发且计数归0时,通知主流程所有任务已处理完毕。
let pendingTaskCount = 0; let isShuttingDown = false; let resolveAllTaskDone = null; const allTaskDonePromise = new Promise(resolve => { resolveAllTaskDone = resolve; }); /** * 追踪异步任务执行状态 * @param {Promise} taskPromise 异步任务对应的Promise * @returns */ function trackTask(taskPromise) { pendingTaskCount++; taskPromise.finally(() => { pendingTaskCount--; // 关闭流程中所有任务执行完成时,触发完成通知 if (isShuttingDown && pendingTaskCount === 0) { resolveAllTaskDone(); } }); return taskPromise; }
2. 所有待保障的异步任务接入追踪
不要保留游离在追踪逻辑外的异步任务,所有需要保证执行完成的逻辑都要通过trackTask包裹:
- HTTP请求:将每个请求的处理逻辑传入
trackTask,替代原有单独的HTTP服务关闭等待逻辑,实现统一追踪 - NATS消息处理:禁止提前ack消息,必须等完整业务逻辑执行完成后再执行ack/nak操作,整个处理流程接入任务追踪。示例订阅逻辑:
const subscription = natsClient.subscribe("your.biz.event", { callback: (err, msg) => { if (err) return; trackTask( (async () => { try { // 执行具体的事件处理逻辑 await handleBizEvent(JSON.parse(msg.data)); // 业务逻辑完全执行完成后再ack msg.ack(); } catch (err) { // 按业务需求执行重试、死信投递、nak等操作 msg.nak({ delay: 5000 }); } })() ); } });
- 自定义后台任务:定时任务、异步批处理等所有后台运行的逻辑,启动时统一接入
trackTask追踪。
3. 按固定顺序编排关闭流程
监听SIGINT、SIGTERM等系统关闭信号,严格按以下顺序执行关闭操作,避免顺序错误导致消息丢失或任务中断:
- 标记服务进入关闭状态,阻止新任务接入
- 关闭HTTP服务,停止接收新的外部HTTP请求
- 调用
nc.drain(),拉取完NATS服务端所有未投递的消息,取消订阅停止接收新的事件消息 - 等待所有追踪中的任务执行完成,增加超时兜底逻辑,避免单个任务卡死导致进程无法退出
- 释放所有外部资源(断开NATS连接、数据库连接、缓存连接等),正常退出进程
关闭逻辑参考代码:
const SHUTDOWN_WAIT_TIMEOUT = 30000; // 最长等待30秒,可按业务调整 async function gracefulShutdown(signal) { // 避免重复触发关闭逻辑 if (isShuttingDown) return; isShuttingDown = true; console.log(`Received ${signal}, start graceful shutdown`); // 关闭HTTP服务,停止接收新请求 await new Promise(resolve => httpServer.close(resolve)); console.log("HTTP server stopped accepting new requests"); // 执行NATS drain,拉取所有未接收消息,停止接新事件 await natsClient.drain(); console.log("NATS client drained, no new events will be received"); // 等待所有在处理任务完成,加超时兜底 const timeoutId = setTimeout(() => { console.error(`Shutdown timeout, force exit, pending tasks: ${pendingTaskCount}`); process.exit(1); }, SHUTDOWN_WAIT_TIMEOUT); if (pendingTaskCount > 0) { console.log(`Waiting for ${pendingTaskCount} running tasks to finish`); await allTaskDonePromise; } clearTimeout(timeoutId); console.log("All pending tasks finished"); // 释放剩余资源 await natsClient.close(); // 此处补充其他资源释放逻辑:断开DB、Redis连接等 console.log("Service shutdown successfully"); process.exit(0); } // 注册关闭信号监听 process.on("SIGINT", () => gracefulShutdown("SIGINT")); process.on("SIGTERM", () => gracefulShutdown("SIGTERM"));
注意事项
- 超时时间需要根据业务实际的最长任务执行时间调整,不要设置过短导致长任务被强制中断,也不要设置过长导致发布流程卡住
- 如果使用NATS JetStream模式,必须严格遵守「业务处理完成再ack」的原则,否则消息会被标记为已消费但实际未处理完成,服务重启后也无法重新消费
- 所有异步任务必须显式接入追踪,不要存在未被
await、未被trackTask包裹的游离异步调用,否则会出现统计遗漏,任务仍会被强制中断
内容的提问来源于stack exchange,提问作者Guna Shekar
相关产品推荐
相关产品推荐

