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

如何获取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,但可以通过以下方式实现:

  1. 判断是否为重发消息:使用msg.getRedelivered()方法,返回布尔值,表示当前消息是否是重发的。
  2. 自定义跟踪投递次数:结合消息的唯一标识(可以是自定义的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 23:48:39