ESP32客户端MQTT消息接收中断,Broker提示消息丢弃问题排查
我使用运行umqtt.simple的ESP32作为客户端,树莓派上运行Paho.MQTT作为Broker,Broker向客户端发送消息,客户端仅负责订阅并在client.call_back()中调用相关函数。我通过client.ping()维持连接存活,初期一切正常,但一段时间后客户端停止接收消息。查看Broker的/var/log/mosquitto/mosquitto.log日志,发现提示Outgoing messages are being dropped for client [CLIENT NAME],且无断开连接记录。
我查资料看到建议在mosquitto.conf中设置max_queued_messages 0,该参数控制每个客户端QoS1或2消息的最大队列数,但我发布消息时使用默认qos=0(即“发后即忘”),理论上不应有消息被缓存。
Broker端简化代码
import paho.mqtt.client as paho BROKER = [broker_ip] uname = [uname] pwd = [password] def on_pub(client, userdata, result): print('Published', result) pass client = paho.Client('Pi') client.on_publish = on_pub client.username_pw_set(password = pwd, username = uname) while True: [满足特定条件时]: # 该条件不会频繁触发,为运动传感器逻辑:一段时间无运动或运动重置计数器到期时发送消息 client.publish(topic = [topic], payload = [payload])
ESP32端MicroPython简化代码
from umqtt.simple import MQTTClient import time [调用WiFi连接函数] BROKER_IP = [broker_IP] CLIENT_NAME = [client_name] USER = [user] PASSWORD = [password] TOPIC = [topic] x = [some_list] def sub_cb(topic, msg): global x [清洗消息字符串转为列表] [将清洗后的列表和x传入某个函数] [重新赋值变量,让清洗后的列表在消息接收期间成为新的x] def connect_and_subscribe(): global CLIENT_NAME, BROKER_IP, USER, PASSWORD, TOPIC client = MQTTClient(client_id=CLIENT_NAME, server=BROKER_IP, user=USER, password=PASSWORD, keepalive=60) client.set_callback(sub_cb) client.connect() client.subscribe(TOPIC) print('Connected to MQTT broker at: %s, subscribed to topic: %s' % (BROKER_IP, TOPIC)) return(client) client = connect_and_subscribe() now = time.time() while True: try: # 防止Broker重启导致check_msg失败,失败后等待60秒重新连接 client.check_msg() except: time.sleep(60) client = connect_and_subscribe() time.sleep(0.1) if time.time() - now > 80: # 每80秒Ping一次维持连接(Broker会在90秒后断开无活动客户端,即1.5倍keepalive值) client.ping() now = time.time()
补充说明:sub_cb中调用的函数可能耗时5秒。我了解Golang的MQTT文档建议为长耗时回调单独开协程,但Python文档未提及相关内容。
我想知道:出现该错误的原因是什么?如何解决?我不想直接设置max_queued_messages 0,担心会导致消息无限缓存直至内存耗尽崩溃,但可能我的理解有误。
原因分析
- 客户端消息处理阻塞:
umqtt.simple是单线程同步库,check_msg()调用后如果sub_cb里的耗时函数阻塞5秒,整个主线程会被卡住,无法处理Broker发来的后续消息,也无法及时响应Broker的网络交互。Broker会认为客户端虽处于连接状态,但无法处理消息,因此开始丢弃待发送的消息。 - Keepalive逻辑的冲突:设置的keepalive为60秒,但你每80秒才发一次PING,加上回调阻塞的5秒,可能导致Broker无法及时收到客户端的存活信号。即使客户端主动发了PING,若消息处理阻塞导致网络IO停滞,Broker仍会判定客户端处理能力不足,触发消息丢弃。
- QoS0的队列误解:虽然QoS0不需要ACK,但Broker会在客户端接收窗口被占满时(即客户端无法及时取走消息),临时缓存QoS0消息,当缓存达到默认上限时就会触发丢弃。
解决方案
1. 异步处理长耗时回调
利用MicroPython的_thread模块,将耗时操作放到子线程执行,避免阻塞主消息循环:
import _thread def sub_cb(topic, msg): global x # 封装耗时操作到子线程函数 def process_msg(): global x [清洗消息字符串转为列表] [将清洗后的列表和x传入某个函数] [重新赋值变量,让清洗后的列表在消息接收期间成为新的x] # 启动子线程执行 _thread.start_new_thread(process_msg, ())
这样sub_cb会立即返回,主循环能持续处理check_msg()和PING操作,客户端可及时响应Broker,避免消息被丢弃。
2. 优化主循环逻辑
调整主循环结构,确保check_msg()被高频调用,同时避免长时间阻塞。比如将耗时操作拆分为多个小步骤,分批次执行,减少单次阻塞时长。
3. 调整Keepalive与PING时机
将PING间隔调整为小于keepalive值的一半(例如keepalive=60时,每25秒发一次PING),确保即使有短暂阻塞,Broker也能及时收到客户端的存活信号。同时保证time.sleep(0.1)不会被阻塞,维持check_msg()的调用频率。
4. 谨慎调整Broker队列参数
若上述方案无法解决,可临时设置max_queued_messages 0允许无限队列,但需结合客户端阻塞问题的修复,否则会导致Broker内存占用上升。实际上,只要客户端恢复正常处理能力,Broker会将队列中的消息逐步发送,不会引发崩溃。
内容的提问来源于stack exchange,提问作者MorrisseyJ

