如何利用Python Paho获取Mosquitto断开连接期间的所有消息
问题描述
我有一个向Mosquitto代理发布消息的Python客户端,该代理通过桥接与第二个代理相连。我希望编写一个Python订阅脚本,获取客户端两次连接期间发布的所有消息,但目前仅能获取到最后一条消息。
主代理(broker)的mosquitto.conf
pid_file /run/mosquitto/mosquitto.pid persistence true persistence_location /var/lib/mosquitto/ log_type all log_dest file /var/log/mosquitto/mosquitto.log log_timestamp true log_timestamp_format %Y-%m-%dT%H:%M:%S include_dir /etc/mosquitto/conf.d listener 1883 allow_anonymous true persistent_client_expiration 14d connection bridge-rpi address 192.168.1.140:1884 topic # out 2 cleansession true remote_clientid xxxx
桥接代理(bridge)的mosquitto.conf
listener 1884 allow_anonymous true persistent_client_expiration 14d persistence true persistence_file mosquitto.db persistence_location /var/lib/mosquitto log_type all log_dest file /var/log/mosquitto/mosquitto.log log_timestamp true log_timestamp_format %Y-%m-%dT%H:%M:%S
当前订阅代码
#Configuration MQTT MQTT_HOST = "localhost" MQTT_PORT = 1884 MQTT_KEEPALIVE_INTERVAL = 45 MQTT_TOPIC = "cn" #Configuration Influxdb INFLUXDB_ADDRESS = 'localhost' INFLUXDB_USER = 'xxxx' INFLUXDB_PASSWORD = 'xxxx' INFLUXDB_DATABASE = 'xxxx' influxdb_client = InfluxDBClient(INFLUXDB_ADDRESS, 8086, INFLUXDB_USER, INFLUXDB_PASSWORD, database=INFLUXDB_DATABASE) # initialisation du client MQTT mqttc = mqtt.Client("PythonMqttToInfluxDb") #mqttc = mqtt.Client() # connexion au broker MQTT mqttc.connect(MQTT_HOST, MQTT_PORT, MQTT_KEEPALIVE_INTERVAL) def on_connect(client, userdata, flags, rc): """ The callback for when the client connects to the broker.""" print("Connecté au serveur MQTT avec le code de retour "+str(rc)) # Subscribe to a topic client.subscribe(MQTT_TOPIC,qos=2) def on_message(client, userdata, msg): """ The callback for when a PUBLISH message is received from the server.""" print(msg.topic+" "+str(msg.payload)) sensor_data = _parse_mqtt_message(msg.topic, msg.payload.decode('utf-8')) json_body = [ { 'measurement': "cn", 'time': parsedJson["timestamp"], 'tags': { 'site': parsedJson["site"], 'frequence': parsedJson["frequence"], 'info': parsedJson["timestamp"] }, 'fields': { 'cn': parsedJson["cn"] #'time': parsedJson["timestamp"] } } ] print(json_body) influxdb_client.write_points(json_body, time_precision='n') ## MQTT logic - Register callbacks and start MQTT client mqttc.on_connect = on_connect mqttc.on_message = on_message mqttc.loop_forever()
解决方案
要获取离线期间的所有消息,需从客户端和代理配置两方面调整:
1. 订阅客户端配置修改
- 创建MQTT客户端时禁用
clean_session,改为False,让代理为客户端保留离线消息:mqttc = mqtt.Client("PythonMqttToInfluxDb", clean_session=False) - 确认订阅的QoS级别≥1(当前用QoS2符合要求),QoS0消息不会被代理保留。
- 修正代码中变量名错误:
parsedJson应改为sensor_data,否则会抛出未定义异常:sensor_data = _parse_mqtt_message(msg.topic, msg.payload.decode('utf-8')) json_body = [ { 'measurement': "cn", 'time': sensor_data["timestamp"], 'tags': { 'site': sensor_data["site"], 'frequence': sensor_data["frequence"], 'info': sensor_data["timestamp"] }, 'fields': { 'cn': sensor_data["cn"] } } ]
2. 桥接代理配置调整
- 修改主代理的桥接配置,将
cleansession设为false,否则桥接代理不会保留主代理转发的离线消息:connection bridge-rpi address 192.168.1.140:1884 topic # out 2 cleansession false remote_clientid xxxx - 确认桥接代理的
persistence已开启(当前配置为persistence true,无需修改),确保重启代理后保留的消息不丢失。
3. 额外注意事项
- 代理的
persistent_client_expiration设为14天,离线超过14天的消息会被清理,符合需求。 - 确保订阅主题
"cn"与发布端主题完全匹配,避免主题过滤导致消息遗漏。
内容的提问来源于stack exchange,提问作者J.bzc
相关产品推荐
相关产品推荐

