如何获取Node NATS Streaming Server消息投递次数及避免重复消费
Node NATS Streaming:错误场景下避免消息重复接收与投递次数获取
一、避免错误场景下的消息重复接收
当服务处理消息抛出异常时,若不手动确认消息,NATS Streaming会持续重发该消息直到确认;但直接确认又可能导致消息丢失。这里提供几种实用方案:
- 死信队列(DLQ)机制:处理失败时,先将消息转发到专门的死信队列,再确认原消息。这样既不会让原订阅端反复收到错误消息,也能保留故障消息用于后续排查和修复。
- 本地重试策略:对临时可恢复的异常(如数据库连接超时、第三方服务暂时不可用),在本地重试2-3次后,再决定是否送入死信队列。
- 业务允许的情况下直接确认:如果业务可以容忍消息丢失,确认后务必记录详细错误日志,方便后续定位问题。
代码示例
const stan = require('node-nats-streaming').connect('your-cluster-id', 'your-client-id'); stan.on('connect', () => { // 开启持久化订阅与手动确认 const sub = stan.subscribe('target-subject', { durableName: 'persistent-sub', manualAckMode: true }); sub.on('message', (msg) => { try { // 执行业务处理逻辑 handleBusinessLogic(msg.getData()); msg.ack(); // 处理成功后确认消息 } catch (err) { console.error(`消息处理失败 [ID: ${msg.getSequence()}]:`, err); // 将消息转发到死信队列 stan.publish('dlq-subject', msg.getData(), (dlqErr) => { if (!dlqErr) { msg.ack(); // 成功转入DLQ后确认原消息,避免重复接收 } else { console.error('死信队列投递失败,将等待NATS重发:', dlqErr); // 此处不确认,让NATS后续重发 } }); } }); }); function handleBusinessLogic(data) { // 你的业务代码 }
二、获取消息投递次数
NATS Streaming的官方Node客户端没有直接提供获取投递次数的API,但可以通过以下方式实现:
- 判断是否为重发消息:使用
msg.getRedelivered()方法,返回布尔值,表示当前消息是否是重发的。 - 自定义跟踪投递次数:结合消息的唯一标识(可以是自定义的
Nats-Msg-Id头,或消息自带的序列号msg.getSequence()),用缓存(如Redis)或本地存储来维护每个消息的投递次数。每次收到消息时,读取并递增计数。
代码示例
const redis = require('redis'); const client = redis.createClient(); // ... 连接NATS的代码省略 sub.on('message', async (msg) => { const msgUniqueId = msg.getHeaders()?.get('Nats-Msg-Id') || msg.getSequence().toString(); // 从Redis获取当前投递次数,初始为0 let deliveryCount = await client.get(`delivery:${msgUniqueId}`); deliveryCount = deliveryCount ? parseInt(deliveryCount) + 1 : 1; // 更新缓存中的次数(设置1天过期) await client.setEx(`delivery:${msgUniqueId}`, 86400, deliveryCount); console.log(`消息 [${msgUniqueId}] 投递次数: ${deliveryCount}`); console.log(`是否为重发: ${msg.getRedelivered()}`); // 后续处理逻辑... });
注意:如果使用持久化订阅,NATS服务器会自动保留未确认的消息并重发,但投递次数需要自行维护。若服务重启,缓存中的计数会丢失,可结合数据库持久化来解决。
内容的提问来源于stack exchange,提问作者Vaibhav Mittal
相关产品推荐
相关产品推荐

