You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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等系统关闭信号,严格按以下顺序执行关闭操作,避免顺序错误导致消息丢失或任务中断:

  1. 标记服务进入关闭状态,阻止新任务接入
  2. 关闭HTTP服务,停止接收新的外部HTTP请求
  3. 调用nc.drain(),拉取完NATS服务端所有未投递的消息,取消订阅停止接收新的事件消息
  4. 等待所有追踪中的任务执行完成,增加超时兜底逻辑,避免单个任务卡死导致进程无法退出
  5. 释放所有外部资源(断开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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 04:09:24