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

paho mqtt的on_message回调如何处理同时到达的多条消息并跨文件传递

原有代码缺陷说明

原有实现核心问题有两个:

  • 使用无锁全局单值变量存储消息,并发场景下后到的消息会直接覆盖未处理的前序消息
  • 普通可变对象(如列表)的读写操作不是原子操作,高并发场景下即使换用列表存储,依然存在竞态条件导致数据丢失的风险
多消息并发处理最优方案

推荐使用Python标准库自带的queue.Queue作为消息缓冲,该容器天生是线程安全的,内部自动实现了锁机制,任意并发量级下都不会出现消息丢失或者读写异常的问题。
改造后的示例代码如下:

import paho.mqtt.client as mqtt
import queue
import threading

# 初始化线程安全队列,maxsize=0代表队列长度无上限,可根据业务需求设置上限避免内存溢出
mqtt_msg_queue = queue.Queue(maxsize=0)

def on_connect(client, userdata, flags, rc):
    print("Connected with result code "+str(rc))
    client.subscribe(topic)

def on_message(client, userdata, message):
    # 直接将消息推入队列即可,无需全局变量存储
    print("received data is :")  
    mqtt_msg_queue.put(message.payload)

# 单独启动工作线程消费消息,避免阻塞MQTT的网络循环线程
def msg_worker():
    while True:
        # 阻塞等待获取消息
        msg = mqtt_msg_queue.get()
        # 此处写你的消息处理逻辑
        print(f"处理消息: {msg}")
        # 标记消息处理完成
        mqtt_msg_queue.task_done()

# 启动消费线程
threading.Thread(target=msg_worker, daemon=True).start()

client = mqtt.Client("user") 
client.on_connect=on_connect
client.on_message=on_message
client.connect(broker,port,60)
client.loop_start()

如果需要分类处理不同topic的消息,可以维护一个topic和队列/处理函数的映射字典,on_message收到消息后根据topic路由到对应处理逻辑即可。注意不要在on_message回调中执行耗时操作,否则会阻塞MQTT的网络线程,导致断连或消息堆积。

跨文件传递消息实现方法

两种常用的低耦合实现方案:

方案1:队列单例模块

单独创建一个全局单例的队列模块,所有文件共享同一个队列实例,完全避免全局变量污染:

  1. 新建mqtt_queue.py文件,内容如下:
import queue
# 全局唯一的消息队列实例
mqtt_msg_queue = queue.Queue(maxsize=0)
  1. MQTT连接文件、业务处理文件都导入该队列实例:MQTT接收端往队列中put消息,其他业务文件从队列中get消息即可。

方案2:回调注册机制

封装MQTT客户端时预留回调注册接口,其他需要接收消息的模块提前将自己的处理函数注册到MQTT客户端,收到消息时自动触发所有注册的回调:

# MQTT客户端模块中维护回调列表
msg_callbacks = []

def register_msg_callback(callback):
    msg_callbacks.append(callback)

def on_message(client, userdata, message):
    for cb in msg_callbacks:
        # 依次调用所有注册的回调函数传递消息
        cb(message.payload)

其他模块只需调用register_msg_callback传入自己的处理函数,即可实时接收MQTT消息,适合多模块同时消费同一条消息的场景。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 04:54:02