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

如何利用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 09:10:35