多MQTT主题订阅失败(连接中断)及重复发布旧消息问题求助
核心问题原因分析
1. 通配符订阅触发自消息循环
代码使用client.subscribe("#", { qos: 1 })订阅所有主题,当程序发布消息到xxx/R时,会被自身的message回调接收,触发重复处理逻辑,最终导致消息循环、连接异常甚至程序退出。
2. 未声明的全局变量引发状态混乱
temp_packet、great_check、count_temp等变量未用var/let/const声明,会挂载到全局作用域。多主题或多消息并行处理时,这些变量的状态会被互相覆盖,导致丢包检测逻辑错误,重复发布旧消息。
3. 定时器重复创建与清理不严谨
每次检测到丢包时都会创建新的setInterval,未检查是否已有同类型定时器在运行;同时setTimeout清理定时器的逻辑仅执行一次,可能导致多个定时器同时发布相同消息。
4. 数据查询空值未处理
Data.find、TopicT.find如果未查询到数据,后续访问id[0]._id、loc[0].last_latitude会直接抛出错误,触发error回调中的process.exit(1),导致程序退出。
5. Async/Await使用不合理
对datTimes.setHours等同步方法使用await无意义且会增加异步回调堆积;message回调是异步函数,但部分数据库操作的错误处理不完整,未捕获异常导致程序崩溃。
修复方案与代码调整
1. 过滤自身发布的消息
在message回调开头添加判断,跳过自己发布的/R主题消息:
client.on('message', async function (topic, message, packet) { // 跳过自身发布的回复主题,避免循环处理 if (topic.endsWith('/R')) { return; } // 原有业务逻辑... });
2. 修复全局变量问题
将所有未声明的变量改为局部变量,用let/const声明;如果需要按设备维护丢包状态,用对象存储:
// 替换原有全局变量 let count = 0; const alphas = ["0", "1", "2", "3", "4", "5", "6", "7", "8", "9", "A", "B", "C", "D", "E", "F", "G", "H", "I", "J", "K", "L", "M", "N", "O", "P", "Q", "R", "S", "T", "U", "V", "W", "X", "Y", "Z"]; // 按设备ID维护上一个数据包编号,避免多设备状态干扰 const devicePacketStates = {};
3. 优化定时器逻辑
添加标志位避免重复创建定时器,同时调整QoS确保消息送达:
// 按设备维护发布状态 const publishingStates = {}; // 丢包检测后的发布逻辑调整 if (value[0] == 1 && !publishingStates[topic_arr[0]]) { publishingStates[topic_arr[0]] = true; const pub_topic = topic_arr[0] + '/R'; const pubmessage = `1,${value[1]}-${value[2]};`; const pub_options = { qos: 1, retain: false }; // 用QoS1保证消息可靠送达 const publishInterval = setInterval(() => { client.publish(pub_topic, pubmessage, pub_options, (err) => { if (err) { console.log("发布错误:", err); } else { console.log("发布成功:", pubmessage); } }); }, 10000); setTimeout(() => { clearInterval(publishInterval); publishingStates[topic_arr[0]] = false; }, 30000); // 按需调整停止发布的时间 }
4. 处理数据查询空值
在使用查询结果前先判断是否存在,避免空值报错:
const id = await Data.find({ id: topic_arr[0] }).catch(err => console.log(err)); if (!id || id.length === 0) { console.log("未找到对应设备数据:", topic_arr[0]); return; } // TopicT查询空值处理 const loc = await TopicT.find({ id: ObjectId(id[0]._id) }, { last_latitude: 1, last_longitude: 1 }); if (!loc || loc.length === 0) { console.log("未找到设备位置数据:", id[0]._id); latitude = message_arr[7]; longitude = message_arr[8]; } else { latitude = Number(message_arr[7]) ? message_arr[7] : loc[0].last_latitude; longitude = Number(message_arr[8]) ? message_arr[8] : loc[0].last_longitude; }
5. 修正Async/Await使用
移除同步方法的await,用try/catch包裹异步逻辑避免未捕获异常:
// 日期处理逻辑简化,移除不必要的await const interMedDt = message_arr[1].replace(/:/gi, " ").replace(/(\w+)\s(\w+)/, "$2 $1"); const datTimes = new Date(interMedDt); const [hours, minutes, seconds] = message_arr[2].split()[0].split(':'); datTimes.setHours(hours); datTimes.setMinutes(minutes); datTimes.setSeconds(seconds); // 用try/catch包裹异步处理逻辑 client.on('message', async function (topic, message, packet) { try { // 原有业务逻辑... } catch (err) { console.log("消息处理错误:", err); // 避免直接退出程序,记录错误后继续运行 } });
6. 调整订阅策略(可选)
如果不需要监听所有主题,明确订阅需要的主题前缀,减少回调压力:
client.subscribe(["+/T", "+/K", "+/S"], { qos: 1 });
内容的提问来源于stack exchange,提问作者Ashuthosh

