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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 23:40:38