Angular使用ngx-mqtt无法正常订阅多个MQTT主题问题咨询
问题根因
你遇到的主题拼接问题核心是两个错误导致的:
- 调用
connect方法后立刻执行订阅,未等待MQTT连接完全建立成功,ngx-mqtt在未就绪状态下会暂存多次订阅请求,内部逻辑异常导致主题被拼接 - 方案1的写法存在逻辑错误,两次对
obs1$赋值会覆盖第一个主题的Observable,残留的未处理订阅请求会和后续订阅逻辑冲突
正确实现方案
首先需要监听MQTT连接成功的事件,等连接就绪后再执行订阅操作,订阅多主题可以直接传入主题数组,也可以分开订阅,两种方式均可正常运行:
import { take } from 'rxjs'; mqttServiceOpts1: IMqttServiceOptions = { connectOnCreate: false, hostname: 'example-mqtt.ca', port: 8090, path: '/mqtt', protocol: 'wss' } connetMqtt() { return new Promise((resolve, reject) => { this.getmqttDetails() .subscribe((data) => { this.mqttServiceOpts1.clientId = data.clientId // 先监听连接成功事件 this.mqttService.onConnect .pipe(take(1)) // 只取一次连接成功事件,避免重复触发 .subscribe(() => { const responsePustatusName= data.rootTopic + '/pustatus/inbox/+/response'; const vacateRooutetopicName= data.rootTopic + '/vacateroute/wc/+/route/response'; // 方式1:一次性订阅多个主题,返回统一的消息流 this.mqttService.observe([responsePustatusName, vacateRooutetopicName]) .subscribe((message: IMqttMessage) => { console.log('收到消息: ', message.payload.toString()) }) // 方式2:分开订阅两个主题,各自处理消息流 // this.mqttService.observe(responsePustatusName) // .subscribe((message: IMqttMessage) => { // console.log('状态主题消息: ', message.payload.toString()) // }) // this.mqttService.observe(vacateRooutetopicName) // .subscribe((message: IMqttMessage) => { // console.log('路线主题消息: ', message.payload.toString()) // }) resolve(true) }) // 监听连接事件后再发起连接 this.mqttService.connect(this.mqttServiceOpts1); }) }) }
额外注意事项
- 不需要自行维护多个Observable变量,ngx-mqtt内部会管理所有主题的订阅关系
- 组件销毁时需要调用
unsubscribe取消所有消息订阅,同时调用disconnect断开MQTT连接,避免内存泄漏
内容的提问来源于stack exchange,提问作者Inder R Singh
相关产品推荐
相关产品推荐

