Django Channels集成MQTT时consumer内on_message回调不触发如何解决
问题解决方法
问题根源
你的代码无法触发updater函数是几个核心问题共同导致的:
- paho-mqtt的
on_message是同步回调接口,你传入异步函数后,paho不会自动执行await操作,函数直接静默失效 - 你将MQTT客户端对象定义为
websocket_connect方法内的局部变量,方法执行完成后对象会被Python垃圾回收,MQTT连接及回调绑定直接失效 - paho的
loop_start会启动独立线程运行MQTT监听逻辑,回调触发时处于MQTT线程中,直接调用Channels的异步send方法属于跨线程调用异步接口,不会正常执行 - 额外问题:你直接将MQTTMessage对象作为text发送给WebSocket,不符合数据格式要求,需要提取payload并解码为字符串
直接修复现有代码的方案
修改后的consumer.py代码如下:
import asyncio from channels.consumer import AsyncConsumer from paho.mqtt import client as Mqtt class Testing(AsyncConsumer): async def websocket_connect(self, event): # 把MQTT客户端存为实例属性,避免被GC回收 self.mqtt_client = Mqtt.Client() self.mqtt_client.connect("localhost", 1883, 60) # 保存当前消费者所在的asyncio事件循环,用于跨线程调度异步方法 self.loop = asyncio.get_running_loop() # 绑定同步回调函数 self.mqtt_client.on_message = self.sync_updater self.mqtt_client.subscribe("Testing") self.mqtt_client.loop_start() # 接受WebSocket连接 await self.send({ "type": "websocket.accept" }) # 同步回调,运行在MQTT线程中 def sync_updater(self, arg1, arg2, message): # 把异步方法调用丢给消费者所在的事件循环执行 asyncio.run_coroutine_threadsafe(self.updater(message), self.loop) async def updater(self, message): # 提取MQTT消息payload解码为字符串 msg_content = message.payload.decode() print(msg_content) await self.send({ "type": "websocket.send", "text": msg_content }) async def websocket_receive(self, text_data): pass # 处理WebSocket断开,清理MQTT资源 async def websocket_disconnect(self, event): self.mqtt_client.loop_stop() self.mqtt_client.disconnect()
更推荐的生产级架构
上述方案每个WebSocket连接都会创建一个独立的MQTT连接,并发高时会导致MQTT服务端连接数超限,推荐用以下架构优化:
- 单独启动一个全局的MQTT客户端服务,仅维持1个MQTT连接
- 收到MQTT消息后通过Channels的Channel Layer发送到指定频道组
- 所有消费者在WebSocket连接建立时加入对应频道组,断开时退出,直接接收频道组内的消息转发给前端,避免重复创建MQTT连接
内容的提问来源于stack exchange,提问作者user12797847
相关产品推荐
相关产品推荐

