调用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
相关产品推荐
相关产品推荐

