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

如何让NATS订阅持续监听并处理消息?

NATS订阅中for await循环内调用异步逻辑后循环退出的解决方法

问题描述

使用nats@2.15.1的发布订阅示例代码时遇到两个问题:

  1. 在for await...of循环内添加异步函数处理消息时,循环会直接退出,无法接收后续新消息;
  2. 不添加异步函数时,订阅虽能持续保持,但console.log始终无法执行。

原示例代码:

(async () => {
  for await (const m of sub) {
    console.log(`[${sub.getProcessed()}]: ${sc.decode(m.data)}`);
  }
  console.log("subscription closed");
})();

解决方案

核心问题分析

  • 循环退出原因:异步处理过程中抛出未捕获的异常,导致for await...of循环终止;或误调用了sub.unsubscribe()等终止订阅的方法。
  • console.log不执行原因:大概率是订阅启用了手动确认模式(manualAck: true)但未确认消息,导致消息堆积阻塞后续接收;或控制台输出被缓冲未刷新。

正确实现代码

要在循环内安全执行异步逻辑且保持订阅活跃,需做好异常捕获与消息确认:

(async () => {
  for await (const m of sub) {
    try {
      // 先打印消息确认接收
      console.log(`[${sub.getProcessed()}]: ${sc.decode(m.data)}`);
      
      // 执行你的异步处理逻辑,比如数据库操作、API调用等
      await yourAsyncMessageHandler(m);
      
      // 若订阅启用了手动确认,必须调用ack()完成消息确认
      // m.ack();
      
      // 可选:处理失败时调用nak()让NATS重新投递消息
      // m.nak();
    } catch (err) {
      // 捕获所有异常,避免终止循环
      console.error("消息处理失败:", err);
    }
  }
  console.log("subscription closed");
})();

关键注意事项

  • 强制异常捕获:所有异步操作必须包裹在try/catch块中,未捕获的异常会直接终止for await...of循环,进而关闭订阅。
  • 消息确认机制:如果创建订阅时设置了{ manualAck: true },处理完成后必须调用m.ack(),否则NATS会认为消息未处理完成,后续消息可能被阻塞。
  • 避免长时间阻塞:若异步逻辑耗时较长,可考虑使用任务池并行处理,防止单个消息处理阻塞后续消息接收;也可改用NATS提供的sub.process()方法实现批量/并行处理。

针对console.log不执行的额外排查点

  1. 确认订阅的主题有消息发布,可手动发送测试消息验证;
  2. 检查订阅配置是否开启了手动确认,若开启需补充m.ack()调用;
  3. Node.js环境下可尝试添加process.stdout.flush()强制刷新控制台输出。

内容的提问来源于stack exchange,提问作者Jason Leach

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 20:08:19