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

多MQTT主题订阅失败(连接中断)及重复发布旧消息问题求助

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 14:30:38