Django中实现MQTT异步Worker:复用连接的方案问询
Django中复用MQTT连接的实现方案
问题背景
在Django应用中,需要从多个位置(包括Celery任务、其他业务逻辑模块)连接MQTT Broker。当前每个场景都要单独创建MQTT客户端并建立连接,不仅重复造轮子,还会因每次等待连接拖慢执行效率,同时缺少连接断开后的自动重连机制,可靠性不足。
期望实现:
- 一个后台/独立线程运行的MQTT Worker,启动时建立连接并持续维护,支持自动重连
- 提供全局可调用的发布、订阅接口,让业务代码只需简单调用即可完成MQTT操作(如示例中的
mqtt_worker.publish())
当前实现痛点
@shared_task( bind=True, name="tasks.send_command", ) def send_command(self, username): pedestals = User.objects.filter(username=username) client = MqttClient("scheduled") client.connect() ... @shared_task( bind=True, name="tasks.toggle_switch", ) def toggle_switch(self, switch): from django.core.cache import cache client = MqttClient("toggle-switch") client.connect() ...
- 每个任务/函数都要重复创建客户端、建立连接
- 每次执行都要等待连接建立,影响效率
- 无自动重连机制,连接断开后业务会直接失效
可行解决方案
方案一:单例模式+后台线程实现独立MQTT Worker
核心思路是用单例保证全局只有一个MQTT客户端实例,在Django启动时初始化并启动后台线程维护连接、处理订阅消息,同时提供简洁的发布接口。
代码实现
创建mqtt/worker.py:
import threading import time import paho.mqtt.client as mqtt from django.conf import settings class MQTTWorker: _instance = None _lock = threading.Lock() def __new__(cls): with cls._lock: if cls._instance is None: cls._instance = super().__new__(cls) cls._instance._client = None cls._instance._connected = False cls._instance._connect_thread = None cls._instance._setup_client() cls._instance._start_connect_loop() return cls._instance def _setup_client(self): self._client = mqtt.Client(client_id=settings.MQTT_CLIENT_ID) self._client.username_pw_set(settings.MQTT_USER, settings.MQTT_PASSWORD) # 设置回调函数 self._client.on_connect = self._on_connect self._client.on_disconnect = self._on_disconnect self._client.on_message = self._on_message def _on_connect(self, client, userdata, flags, rc): if rc == 0: self._connected = True # 订阅业务所需主题 self._client.subscribe(settings.MQTT_SUBSCRIBE_TOPICS) else: print(f"MQTT连接失败,错误码: {rc}") def _on_disconnect(self, client, userdata, rc): self._connected = False print("MQTT连接断开,将尝试重连") def _on_message(self, client, userdata, msg): # 处理订阅到的消息,可根据业务需求分发(如触发Django信号、调用Celery任务) print(f"收到MQTT消息: 主题={msg.topic}, 内容={msg.payload.decode()}") # 示例:触发信号处理消息 # from mqtt.signals import mqtt_message_received # mqtt_message_received.send(sender=self, topic=msg.topic, payload=msg.payload) def _connect_loop(self): while True: if not self._connected: try: self._client.connect(settings.MQTT_BROKER_HOST, settings.MQTT_BROKER_PORT, keepalive=60) self._client.loop_start() # 启动非阻塞循环 except Exception as e: print(f"MQTT连接出错: {e}") time.sleep(5) # 重连间隔 time.sleep(1) def _start_connect_loop(self): self._connect_thread = threading.Thread(target=self._connect_loop, daemon=True) self._connect_thread.start() def publish(self, topic, payload, qos=0, retain=False): if self._connected: self._client.publish(topic, payload, qos, retain) else: print("MQTT未连接,无法发布消息") # 全局可调用的实例 mqtt_worker = MQTTWorker()
在Django的mqtt/apps.py中配置启动加载:
from django.apps import AppConfig class MqttConfig(AppConfig): default_auto_field = 'django.db.models.BigAutoField' name = 'mqtt' def ready(self): # 初始化MQTT Worker,确保Django启动时加载 from .worker import mqtt_worker
业务代码中使用
@shared_task(bind=True, name="tasks.toggle_switch") def toggle_switch(self, switch): from django.core.cache import cache from mqtt.worker import mqtt_worker # 直接调用发布接口 mqtt_worker.publish(topic="device/switch", payload=f"toggle:{switch}") ...
方案二:结合Celery实现MQTT Worker
利用Celery的长期运行任务能力维护MQTT连接,借助Celery的进程管理机制避免连接意外终止。
代码实现
创建mqtt/tasks.py:
from celery import shared_task import paho.mqtt.client as mqtt from django.conf import settings import time mqtt_client = None connected = False def on_connect(client, userdata, flags, rc): global connected if rc == 0: connected = True client.subscribe(settings.MQTT_SUBSCRIBE_TOPICS) else: print(f"MQTT连接失败: {rc}") def on_disconnect(client, userdata, rc): global connected connected = False print("MQTT断开连接,准备重连") def on_message(client, userdata, msg): # 处理订阅消息,可调用其他Celery任务异步处理 print(f"收到消息: {msg.topic} -> {msg.payload}") # from tasks import handle_mqtt_message # handle_mqtt_message.delay(msg.topic, msg.payload.decode()) @shared_task(bind=True, name="mqtt.worker") def mqtt_worker_task(self): global mqtt_client, connected mqtt_client = mqtt.Client(client_id=settings.MQTT_CELERY_CLIENT_ID) mqtt_client.username_pw_set(settings.MQTT_USER, settings.MQTT_PASSWORD) mqtt_client.on_connect = on_connect mqtt_client.on_disconnect = on_disconnect mqtt_client.on_message = on_message while True: if not connected: try: mqtt_client.connect(settings.MQTT_BROKER_HOST, settings.MQTT_BROKER_PORT, 60) mqtt_client.loop_start() except Exception as e: print(f"连接出错: {e}") time.sleep(5) time.sleep(1) # 全局发布工具函数 def mqtt_publish(topic, payload, qos=0, retain=False): global connected, mqtt_client if connected and mqtt_client: mqtt_client.publish(topic, payload, qos, retain) else: print("MQTT未连接,无法发布")
启动Celery Worker并指定该任务单独运行:
celery -A your_project_name worker -l info -Q mqtt_worker --concurrency=1
业务代码中使用
@shared_task(bind=True, name="tasks.toggle_switch") def toggle_switch(self, switch): from django.core.cache import cache from mqtt.tasks import mqtt_publish mqtt_publish(topic="device/switch", payload=f"toggle:{switch}") ...
mqttasgi适配分析
mqttasgi主要为ASGI应用(如Django Channels)设计,适合需要Web端与MQTT双向通信的场景:
- 可自动维护MQTT连接,支持订阅和发布
- 能通过Channels消费者处理订阅消息,无缝集成Django异步生态
- 但如果你的场景以后台任务、同步业务逻辑为主,前面的单例/ Celery长期任务方案更轻量直接
内容的提问来源于stack exchange,提问作者Kristof Rado
相关产品推荐
相关产品推荐

