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

ESP32客户端MQTT消息接收中断,Broker提示消息丢弃问题排查

问题: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,担心会导致消息无限缓存直至内存耗尽崩溃,但可能我的理解有误。


原因分析

  1. 客户端消息处理阻塞:umqtt.simple是单线程同步库,check_msg()调用后如果sub_cb里的耗时函数阻塞5秒,整个主线程会被卡住,无法处理Broker发来的后续消息,也无法及时响应Broker的网络交互。Broker会认为客户端虽处于连接状态,但无法处理消息,因此开始丢弃待发送的消息。
  2. Keepalive逻辑的冲突:设置的keepalive为60秒,但你每80秒才发一次PING,加上回调阻塞的5秒,可能导致Broker无法及时收到客户端的存活信号。即使客户端主动发了PING,若消息处理阻塞导致网络IO停滞,Broker仍会判定客户端处理能力不足,触发消息丢弃。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 22:15:27