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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 21:12:18