MQTT主题数据无法持久化、初始不可见问题技术咨询
问题排查与解决方案
一、现有发布代码的问题
未等待消息发布完成就断开连接
paho.mqtt.client的publish是异步操作,调用后不会立即完成消息发送。你的代码在调用publish后立刻执行disconnect,此时消息可能还在发送队列中就被中断,导致Broker根本没收到消息。
解决:需要等待on_publish回调触发后再断开连接,或者使用同步阻塞的方式确认发送。多余的订阅操作
发布程序不需要订阅目标主题,mqttc.subscribe(MQTT_TOPIC)属于无效操作,可直接删除。未设置消息保留(Retain)与QoS
默认情况下,MQTT消息是临时的,只有当前在线的订阅者能收到。如果需要让后续订阅的客户端(比如MQTT Explorer打开时)能拿到历史最新消息,必须设置retain=True;同时QoS级别如果设为0,Broker不会持久化消息,建议根据需求设置QoS 1或2。
修正后的发布代码:
import paho.mqtt.client as mqtt import json import os import time MQTT_HOST = os.environ["MQTT_HOST"] MQTT_PORT = int(os.environ["MQTT_PORT"]) MQTT_KEEPALIVE_INTERVAL = int(os.environ["MQTT_KEEPALIVE_INTERVAL"]) MQTT_TOPIC = os.environ["MQTT_TOPIC"] # 标记消息是否发布完成 published_flag = False def on_publish(client, userdata, mid): global published_flag published_flag = True print("Message Published...") def on_connect(client, userdata, flags, rc): if rc == 0: print("Connected to MQTT Broker successfully") else: print(f"Failed to connect, return code {rc}") mqttc = mqtt.Client() mqttc.on_publish = on_publish mqttc.on_connect = on_connect mqttc.connect(MQTT_HOST, MQTT_PORT, MQTT_KEEPALIVE_INTERVAL) mqttc.loop_start() # 等待连接稳定(可选,但避免连接未完成就发布) time.sleep(1) print("publishing ") # 设置QoS=1,retain=True,确保消息被Broker持久化并能被新订阅者获取 mqttc.publish(MQTT_TOPIC, json.dumps(reading_table_data), qos=1, retain=True) # 等待消息发布完成 while not published_flag: time.sleep(0.1) mqttc.disconnect() mqttc.loop_stop()
二、消息无法持久化的原因排查
Broker配置问题
确认你的MQTT Broker(如Mosquitto)是否开启了消息持久化功能:- Mosquitto需在配置文件中设置
persistence true,并指定persistence_location存储持久化数据。 - 如果Broker是临时部署(比如K8s中无持久化存储),重启后所有持久化数据会丢失。
- Mosquitto需在配置文件中设置
消息属性问题
- 必须在发布时设置
retain=True,Broker才会保留该主题的最后一条消息,供新订阅者获取。 - QoS 0的消息Broker不会持久化,建议使用QoS 1或2,确保消息至少被送达一次且Broker会持久化未确认的消息。
- 必须在发布时设置
MQTT Explorer显示问题
确保MQTT Explorer订阅主题时勾选了“订阅保留消息”选项,否则即使Broker有保留消息,也不会主动拉取。
三、订阅主题并解析JSON的示例代码
以下是使用paho.mqtt.client订阅目标主题并解析JSON字段的示例:
import paho.mqtt.client as mqtt import json import os MQTT_HOST = os.environ["MQTT_HOST"] MQTT_PORT = int(os.environ["MQTT_PORT"]) MQTT_KEEPALIVE_INTERVAL = int(os.environ["MQTT_KEEPALIVE_INTERVAL"]) MQTT_TOPIC = os.environ["MQTT_TOPIC"] def on_connect(client, userdata, flags, rc): print(f"Connected with result code {rc}") # 订阅目标主题,QoS与发布端保持一致 client.subscribe(MQTT_TOPIC, qos=1) def on_message(client, userdata, msg): print(f"Received message on topic {msg.topic}") try: # 解析JSON数据 data = json.loads(msg.payload.decode()) # 提取需要的字段 energy_source = data.get("Energy Source") grid_reading = data.get("Grid Reading ") power_factor = data.get("Power Factor") account_balance = data.get("accountBalance") print(f"Energy Source: {energy_source}") print(f"Grid Reading: {grid_reading}") print(f"Power Factor: {power_factor}") print(f"Account Balance: {account_balance}") except json.JSONDecodeError as e: print(f"Failed to parse JSON: {e}") mqttc = mqtt.Client() mqttc.on_connect = on_connect mqttc.on_message = on_message mqttc.connect(MQTT_HOST, MQTT_PORT, MQTT_KEEPALIVE_INTERVAL) # 持续运行循环处理消息 mqttc.loop_forever()
内容的提问来源于stack exchange,提问作者Rahul Sharma
相关产品推荐
相关产品推荐

