如何让NATS订阅持续监听并处理消息?
NATS订阅中for await循环内调用异步逻辑后循环退出的解决方法
问题描述
使用nats@2.15.1的发布订阅示例代码时遇到两个问题:
- 在
for await...of循环内添加异步函数处理消息时,循环会直接退出,无法接收后续新消息; - 不添加异步函数时,订阅虽能持续保持,但
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不执行的额外排查点
- 确认订阅的主题有消息发布,可手动发送测试消息验证;
- 检查订阅配置是否开启了手动确认,若开启需补充
m.ack()调用; - Node.js环境下可尝试添加
process.stdout.flush()强制刷新控制台输出。
内容的提问来源于stack exchange,提问作者Jason Leach
相关产品推荐
相关产品推荐

