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
相关产品推荐
相关产品推荐

