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

调用API取消NATS Streaming(STAN)订阅报无效订阅错误解决方案

问题根因

取消订阅逻辑存在两个核心错误:

  • 未持久化存储已创建成功的订阅实例,触发取消操作时重新调用subscribe()方法新建了一个完全独立的冗余订阅,操作对象根本不是需要取消的原有订阅
  • 新建订阅后立刻执行unsubscribe(),此时订阅请求尚未完成STAN服务端的注册确认,持有的订阅句柄为无效状态,因此抛出stan: invalid subscription错误。即便等待该新订阅注册完成后再执行取消,也仅会关闭刚创建的冗余订阅,原有订阅仍会正常接收消息,最终表现为取消功能完全不生效。
修复步骤

1. 新增订阅实例存储结构

新增全局缓存结构存储已创建的订阅实例,以「客户端ID+主题+队列名」为唯一key,避免后续无法定位到需要操作的原订阅:

// 全局订阅实例存储,多client场景下key必须带上clientId避免冲突
const subscriptionMap = new Map();

2. 改造订阅方法,创建成功后持久化实例

调整subscribeToTopic逻辑,等待订阅就绪后再绑定事件、存入缓存,同时增加重复订阅校验:

subscribeToTopic: async function(clientId, topic, queueName) {
    const subKey = `${clientId}:${topic}:${queueName}`;
    // 已存在对应订阅直接返回,避免重复创建
    if (subscriptionMap.has(subKey)) {
        console.log(`subscription for ${subKey} already exists`);
        return subscriptionMap.get(subKey);
    }

    console.log('subscriber function start');
    const options = stanSubscribe.subscriptionOptions().setManualAckMode(true);
    const subscribe = stanSubscribe.subscribe(topic, queueName, options);

    // 等待订阅完成服务端注册就绪
    await new Promise((resolve, reject) => {
        subscribe.on('ready', () => {
            console.log(`subscribed success: ${subKey}`);
            // 绑定消息监听
            subscribe.on('message', async (subMessage) => {
                console.log('Recieved the message, queue name = [' + queueName + '] sequence = [' + subMessage.getSequence() + '] - ' + subMessage.getData());
                // 注意:手动Ack模式下消息处理完成必须调用ack(),否则消息会重复投递
                // subMessage.ack();
            });
            // 订阅实例存入缓存
            subscriptionMap.set(subKey, subscribe);
            resolve(subscribe);
        });
        subscribe.on('error', (err) => {
            console.log(`subscribe failed for ${subKey}:`, err);
            reject(err);
        });
    });

    console.log('subscriber function end');
}

3. 重写取消订阅逻辑,操作原有订阅实例

取消订阅时直接从缓存中取出已创建的目标订阅实例执行操作,禁止重新调用subscribe创建新订阅:

async function unsubscribeFromTopic(clientId, topic, queueName) {
    const subKey = `${clientId}:${topic}:${queueName}`;
    const targetSub = subscriptionMap.get(subKey);
    
    if (!targetSub) {
        console.log(`no active subscription found for ${subKey}`);
        return;
    }

    try {
        // 绑定取消成功回调
        targetSub.on('unsubscribed', () => {
            console.log(`unsubscribed from queue: ${queueName} from topic: ${topic}`);
            // 从缓存中移除已取消的订阅
            subscriptionMap.delete(subKey);
        });

        // 对有效原订阅执行取消操作
        targetSub.unsubscribe();

        if (targetSub.isClosed()) {
            console.log(`subscription is closed for queueName: ${queueName} from topic: ${topic}`);
            subscriptionMap.delete(subKey);
        }
    } catch (error) {
        console.log(`unsubscribe error: ${error}`);
    }
}
额外注意事项
  • 开启手动Ack模式(setManualAckMode(true))后,消息业务逻辑处理完成必须调用subMessage.ack()告知STAN消费完成,否则达到超时阈值后消息会被重复投递
  • 多实例部署场景下,不能仅将订阅实例存在服务本地内存,需要保证订阅创建请求和取消请求路由到同一个服务节点,否则会无法读取到对应订阅实例
  • 订阅实例未触发ready事件前,不要执行unsubscribe、close类操作,此时服务端尚未完成订阅注册,操作无效句柄会触发参数错误

内容的提问来源于stack exchange,提问作者Ashwin Pandey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 03:54:33