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

Paho MQTT运行1-2天后停止收发消息问题求助

MQTT客户端运行1-2天后停止收发消息,publish无报错且is_connected返回True

我用Python结合Paho MQTT开发时遇到异常:代码订阅了incoming主题,收到消息就启动独立线程执行阻塞的control_device函数(耗时约5分钟,且确保同一时间只有一个实例运行),同时每60秒向outgoing主题发布状态信息。程序初期运行正常,但1-2天后会突然停止接收和发布MQTT消息,publish无报错,is_connected()返回True,重启后又能正常运行1-2天。

附上完整代码:

import paho.mqtt.client as mqtt
import json
import threading

class Controller():
    def __init__(self):
        self.mqtt_client = mqtt.Client()
        self.pub_topic = "outgoing"
        self.mqtt_client.on_message = self.on_message
        self.mqtt_client.connect("192.168.1.1", 1883, 600)
        self.mqtt_client.subscribe("incoming")

    # 阻塞函数,执行耗时约5分钟,同一时间仅允许一个实例运行
    def control_device(self, input_commands):
        print("Do some stuff...")

    def process_mqtt(self, msg):
        mqtt_msg = json.loads(msg.payload.decode('utf-8'))
        self.control_device(mqtt_msg)
        payload = '{"message":  "process started"}'                                                                                                       
        self.mqtt_client.publish(self.pub_topic, payload)

    def on_message(self, client, userdata, msg):
        thread = threading.Thread(target=self.process_mqtt, args=(msg,))
        thread.start()

    # 每60秒向同一主题发送状态信息
    def send_status_msg(self):
        if minute_passed:
            payload = '{"status": 0}'
            self.mqtt_client.publish(self.pub_topic, payload)

    def run(self):
        while True:
            self.mqtt_client.loop()
            self.send_status_msg()

if __name__ == "__main__":
    c = Controller()
    c.run()

我的疑问:我是否对MQTT库的工作机制存在误解?看到有讨论说不应在on_message()内发布消息,但我已经把相关操作放到单独线程里了。


问题分析与修复方案

1. loop()调用方式是核心问题

当前用while True循环调用mqtt_client.loop()属于单线程阻塞式调用,会导致MQTT客户端的网络IO处理被send_status_msg的高频调用打断。长期运行后,可能出现心跳包处理不及时、网络队列堆积,最终让客户端处于假在线状态——is_connected()返回True,但实际无法与Broker正常收发消息。

Paho MQTT官方推荐两种更可靠的运行方式:

  • 用loop_start()启动后台线程自动处理网络IO,无需手动循环调用loop()
  • 若必须手动控制,改用loop(timeout=1.0)设置超时,给网络处理留出足够时间

2. 状态发布的计时逻辑缺失

代码里minute_passed变量未定义,实际运行会直接抛出异常导致run()循环崩溃。就算你实际代码中有定义,也得确保计时逻辑线程安全且准确,避免频繁发消息占用带宽。

3. 线程管理存在隐患

虽然on_message里启动了新线程处理任务,但没限制线程数量,也没设置守护属性。短时间内收到大量消息的话,会创建过多线程耗尽系统资源,间接影响MQTT客户端运行。

修复后的代码示例

import paho.mqtt.client as mqtt
import json
import threading
import time

class Controller():
    def __init__(self):
        self.mqtt_client = mqtt.Client()
        self.pub_topic = "outgoing"
        self.mqtt_client.on_message = self.on_message
        # 设置自动重连参数,断连后自动重试,避免假在线
        self.mqtt_client.reconnect_delay_set(min_delay=1, max_delay=60)
        self.mqtt_client.connect("192.168.1.1", 1883, 600)
        self.mqtt_client.subscribe("incoming")
        # 启动后台线程处理MQTT网络IO,无需手动调用loop()
        self.mqtt_client.loop_start()
        # 线程锁确保同一时间只有一个control_device实例运行
        self.device_lock = threading.Lock()
        # 记录上次发送状态的时间
        self.last_status_time = time.time()

    def control_device(self, input_commands):
        print("Do some stuff...")
        # 模拟5分钟阻塞操作
        time.sleep(300)

    def process_mqtt(self, msg):
        # 尝试获取锁,非阻塞模式避免等待
        if self.device_lock.acquire(blocking=False):
            try:
                mqtt_msg = json.loads(msg.payload.decode('utf-8'))
                self.control_device(mqtt_msg)
                payload = '{"message":  "process started"}'
                # Paho的publish方法是线程安全的,子线程中调用没问题
                self.mqtt_client.publish(self.pub_topic, payload)
            finally:
                # 无论是否异常都释放锁
                self.device_lock.release()
        else:
            print("设备控制已在运行,跳过本次指令")

    def on_message(self, client, userdata, msg):
        # 设置守护线程,主线程退出时自动销毁子线程
        thread = threading.Thread(target=self.process_mqtt, args=(msg,), daemon=True)
        thread.start()

    def send_status_msg(self):
        current_time = time.time()
        # 每60秒发送一次状态
        if current_time - self.last_status_time >= 60:
            payload = '{"status": 0}'
            self.mqtt_client.publish(self.pub_topic, payload)
            self.last_status_time = current_time

    def run(self):
        while True:
            self.send_status_msg()
            # 每秒检查一次,降低CPU占用
            time.sleep(1)

if __name__ == "__main__":
    c = Controller()
    try:
        c.run()
    except KeyboardInterrupt:
        # 停止MQTT后台线程,优雅退出
        c.mqtt_client.loop_stop()

关键修复点说明

  • 用loop_start()替代手动loop()调用,让MQTT客户端在后台独立处理网络IO,避免主线程阻塞影响心跳和消息收发
  • 添加device_lock线程锁,确保control_device同一时间只运行一个实例(实现你代码注释里的逻辑)
  • 修复状态发布的计时逻辑,确保每60秒发送一次,同时加入time.sleep(1)降低主线程CPU占用
  • 设置子线程为守护线程,避免主线程退出后残留无用线程
  • 配置自动重连参数,增强连接稳定性,应对临时网络波动
  • 关于on_message内发布消息的疑问:你把操作放到子线程的做法是正确的,之前的讨论主要是指不要在on_message回调里做阻塞操作,Paho的publish方法本身是线程安全的,子线程中调用完全没问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 00:01:20