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:队列单例模块
单独创建一个全局单例的队列模块,所有文件共享同一个队列实例,完全避免全局变量污染:
- 新建
mqtt_queue.py文件,内容如下:
import queue # 全局唯一的消息队列实例 mqtt_msg_queue = queue.Queue(maxsize=0)
- 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
相关产品推荐
相关产品推荐

