使用mqttasgi实现MQTT监听经WSS端点推送数据的方法
结论
完全可以在自有服务器上同时实现MQTT服务监听、WSS端点消息推送。mqttasgi本身基于Django Channels的ASGI生态开发,和Channels原生WebSocket消费者共享同一套Channel Layer通信机制,不需要额外引入第三方组件就能打通MQTT消息流和WebSocket推送链路,完全适配你当前的业务场景。
具体实现步骤
你现有MQTT消费者的核心逻辑不需要大改,只需要补全WebSocket相关配置和代码即可:
- 基础依赖检查
确保环境中已经安装channels和daphne(ASGI服务器,可同时承载MQTT和WebSocket请求),Django配置文件中正确配置Channel Layer,生产环境推荐用Redis作为Channel Layer后端,默认内存后端跨进程不生效,不要在生产环境使用。
WSS不需要在业务代码层做特殊适配,生产环境只需要在Daphne前挂Nginx做SSL终止,将wss请求转发到Daphne监听端口即可。 - 编写WebSocket消费者
新建WebSocket消费者,负责和前端建立WSS连接,加入和MQTT消费者相同的通信组,即可收到MQTT链路处理完成的消息,直接推送给前端:from channels.generic.websocket import AsyncWebsocketConsumer import json class MqttStreamConsumer(AsyncWebsocketConsumer): async def connect(self): # 此处可添加鉴权逻辑,比如校验连接token、用户权限,避免接口裸奔 self.group_name = "stracontech" await self.channel_layer.group_add( self.group_name, self.channel_name ) await self.accept() async def disconnect(self, close_code): await self.channel_layer.group_discard( self.group_name, self.channel_name ) async def receive(self, text_data): # 如果需要前端通过WebSocket下发指令控制MQTT发布/订阅,逻辑可写在这里 pass # 方法名需要和Channel Layer发送事件的type字段对应,用于接收组内广播消息 async def publish_results(self, event): data = event['result'] await self.send(text_data=json.dumps({ "topic": f"stracontech/procesed/{data['device_id']}/result", "payload": data })) - 调整Celery任务回调逻辑
你当前用的processmqttmessage是Celery异步任务,任务处理完结果后,调用Channel Layer的group_send方法把结果广播到stracontech组即可,组内的MQTT消费者和WebSocket消费者都会收到事件:
你原有MQTT消费者中的publish_results方法不需要修改,收到事件后仍然照常往MQTT响应主题发消息,新加的WebSocket消费者会同步把处理结果推送给所有已连接的WSS客户端。
Celery任务末尾添加广播逻辑的参考代码:from channels.layers import get_channel_layer from asgiref.sync import async_to_sync channel_layer = get_channel_layer() async_to_sync(channel_layer.group_send)( "stracontech", { "type": "publish.results", "result": 处理完成的结果字典 } ) - 配置ASGI路由
在ASGI路由配置中同时注册MQTT消费者和WebSocket消费者,参考配置:from django.urls import re_path from channels.auth import AuthMiddlewareStack from channels.routing import ProtocolTypeRouter, URLRouter from .consumers import MyMqttConsumer, MqttStreamConsumer application = ProtocolTypeRouter({ "websocket": AuthMiddlewareStack( URLRouter([ re_path(r'ws/mqtt-data-stream/$', MqttStreamConsumer.as_asgi()), ]) ), "mqtt": URLRouter([ MyMqttConsumer.as_asgi(), ]) }) - 生产环境WSS配置
用Nginx反向代理Daphne服务,配置SSL证书,将WSS请求转发到Daphne服务端口即可,核心配置片段:server { listen 443 ssl; server_name 你的服务域名; ssl_certificate 你的SSL证书路径; ssl_certificate_key 你的SSL证书私钥路径; location /ws/mqtt-data-stream/ { proxy_pass http://127.0.0.1:8000; # 替换为你Daphne实际监听的地址 proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_set_header Host $host; proxy_read_timeout 86400; # 调大长连接超时时间,避免连接被主动断开 } }
常见注意点
- Channel Layer是MQTT和WebSocket消费者通信的核心,必须确保Redis服务稳定,不要用内存后端。
- WebSocket连接务必加上鉴权逻辑,避免未授权访问泄露设备上报的敏感数据。
- 如果前端访问存在跨域问题,在Django跨域配置中添加对应的允许规则即可。
内容的提问来源于stack exchange,提问作者Jonathan Prieto
相关产品推荐
相关产品推荐

