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

MQTT Subscriber中loop_forever()失效问题求助

问题分析与解决方案

核心问题

程序运行6-10小时后重连并停止接收数据,根源在于两个关键点:

  1. 重连后订阅关系丢失:Paho MQTT默认clean_session=True,当连接因心跳超时或网络波动断开后重连,Broker会清除该客户端的会话信息,初始的订阅关系直接失效,而原代码仅在首次启动时执行订阅操作。
  2. 无重连后订阅恢复逻辑:即便客户端自动重连成功,也没有触发重新订阅主题的逻辑,导致无法继续接收Broker推送的数据。

修复步骤

1. 在on_connect回调中嵌入订阅逻辑

无论首次连接还是重连成功,都自动执行订阅操作,确保连接恢复后立即恢复数据接收能力。

2. 启用会话持久化(可选但推荐)

设置clean_session=False,让Broker保留客户端的会话信息,重连时自动恢复订阅,同时配合QoS>0还能接收离线期间未送达的消息。

3. 添加断开连接监控(可选)

增加on_disconnect回调,便于排查异常断开的原因。

修改后的完整代码

# python3.6
import random
import mysql.connector
from paho.mqtt import client as mqtt_client
import json

# MQTT Connection Configuration
broker = 'YOUR_BROKER'
port = 1883
topic = "YOUR_TOPIC"
# Generate client ID with pub prefix randomly
client_id = f'python-mqtt-{random.randint(0, 100)}'
username = "THE_USERNAME"
password = "THE_PASSWORD"
keepalive = 60  # Heartbeat interval in seconds


def connect_mqtt() -> mqtt_client:
    def on_connect(client, userdata, flags, rc):
        if rc == 0:
            print("Connected to MQTT Broker!")
            # 首次连接/重连时自动订阅主题
            client.subscribe(topic, qos=1)
            print(f"Subscribed to topic: {topic}")
        else:
            print(f"Failed to connect, return code {rc}")

    def on_disconnect(client, userdata, rc):
        if rc != 0:
            print(f"Unexpected disconnection, return code {rc}. Reconnecting...")

    # 设置clean_session=False保留会话,重连时自动恢复订阅
    client = mqtt_client.Client(client_id, clean_session=False)
    client.username_pw_set(username, password)
    client.on_connect = on_connect
    client.on_disconnect = on_disconnect
    client.connect(broker, port, keepalive=keepalive)
    return client


def setup_message_handler(client: mqtt_client):
    def on_message(client, userdata, msg):
        print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic")

    client.on_message = on_message


def run():
    client = connect_mqtt()
    setup_message_handler(client)
    client.loop_forever()


if __name__ == '__main__':
    run()

额外说明

  • clean_session=False需保证client_id唯一且固定(你的代码中client_id在程序启动时随机生成,只要不重启就固定,符合要求),否则多个客户端使用同一client_id会导致会话冲突。
  • 如果需要接收离线期间的消息,必须将订阅的QoS设置为1或2(示例中用了qos=1),同时Broker需开启消息持久化配置。

内容的提问来源于stack exchange,提问作者Philipp

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 15:35:48