侧边栏壁纸
  • 累计撰写 30 篇文章
  • 累计收到 2 条评论

Python订阅MQTT消息乱序丢包的排障记录

2026-9-16 / 0 评论 / 8 阅读

Python订阅MQTT消息乱序丢包的排障记录

上周三蹲在客厅给新焊的ESP32温湿度节点调上报,本来想着写个10行的Python脚本订阅MQTT消息存到本地SQLite,跑了没半小时就发现数据时间戳不对,有前几分钟的消息晚到,还有整段的记录缺口。

一开始以为是硬件的锅

碰到数据丢了第一反应就是ESP32供电不稳,毕竟之前吃够了劣质USB线导致模块重启的亏。拆外壳换线、换传感器、串口打日志看上报状态,折腾了快四十分钟,确认模块端每两秒一次的上报稳得很,根本没丢包。转身上EMQX后台看消息流转记录,所有消息都正常投递到了我写的订阅客户端,问题出在脚本本身。说白了当时就是图快,觉得逻辑不复杂怎么写都行,完全忘了回调函数的执行上下文是啥,最开始的错版代码长这样:

import paho.mqtt.client as mqtt
import json
def on_message(client, userdata, msg):
    # 所有逻辑全塞在回调里同步跑
    data = json.loads(msg.payload)
    db_insert(data) # 同步写SQLite
    if data['temp'] > 30:
        send_wechat() # 同步调用webhook发通知
client = mqtt.Client()
client.connect("192.168.31.200",1883,60)
client.subscribe("sensors/#",0)
client.loop_forever()

paho-mqtt默认的loop_forever是单线程跑网络循环的,on_message回调就运行在这个核心线程里,但凡数据库锁一下、通知接口响应慢个几百毫秒,后续的消息全堆在socket缓冲区读不出来,等前面的操作跑完才会按接收顺序处理,自然就出现了消息晚到、甚至缓冲区满了丢包的情况。

试错的时候踩了另一个坑

发现问题的时候想当然觉得,把回调里的耗时逻辑扔到线程池里异步跑不就完了?结果改完之后跑起来时不时报随机的连接错误,翻了半小时官方文档才反应过来,paho-mqtt的客户端实例本身不是线程安全的,在非网络循环的线程里随便调用client的方法,很容易把内部的连接状态搞乱。

最后改完能稳定跑的写法

核心思路就是把收消息和处理消息完全拆开,谁也不堵谁,几步就能改完:

  • 不要在on_message回调里写任何业务逻辑,回调只负责把收到的消息存入队列
  • 单独启动工作线程,从队列里取消息做入库、发通知这类耗时操作
  • 用固定大小的队列做缓冲,队列满的时候主动丢最旧的消息,绝对不能让入队操作阻塞回调
  • 用loop_start()启动独立的网络循环线程,别用loop_forever()阻塞主线程

核心代码片段可以直接复用:

import paho.mqtt.client as mqtt
import queue, threading, json, time
msg_queue = queue.Queue(maxsize=100)

def on_message(client, userdata, msg):
    try:
        msg_queue.put_nowait((msg.topic, msg.payload))
    except queue.Full:
        msg_queue.get_nowait()
        msg_queue.put_nowait((msg.topic, msg.payload))

def msg_worker():
    while True:
        topic, payload = msg_queue.get()
        try:
            data = json.loads(payload)
            db_insert(data)
            if data.get('temp',0) > 30:
                send_wechat()
        except Exception as e:
            print(f"process msg error: {e}")
        finally:
            msg_queue.task_done()

if __name__ == "__main__":
    threading.Thread(target=msg_worker, daemon=True).start()
    client = mqtt.Client(client_id="sensor_collector", clean_session=False)
    client.connect("192.168.31.200", 1883, 60)
    client.subscribe("sensors/#", qos=1)
    client.loop_start()
    while True:
        time.sleep(60)
别图省事在on_message里每次收到消息就新开线程处理,消息量大的时候线程数会失控,用固定大小队列+单工作线程的写法稳得多。家用传感场景QoS设1就够用,开QoS2的交互开销反而容易增加延迟。

改完挂在树莓派上跑到现在,大概一周左右,再也没出现过消息乱序晚到的问题。说起来一开始上来就拆硬件查供电的操作,纯纯是把简单问题想复杂了。

评论一下?

OωO
取消