如何在数据库变更时触发运行中的MQTT监控客户端向Mosquitto发消息?
解决Django信号中复用常驻MQTT客户端发送消息的问题
我完全懂你的困扰:每次用publish.single都会新建一个MQTT连接,既浪费资源,又没法复用你已经跑起来的monitor.py客户端。下面是一个实用的解决方案,核心思路是用**进程间通信(IPC)**让Django把待发送的消息传递给常驻的monitor客户端,由它来完成MQTT发送。
方案概述
因为monitor.py是独立的常驻进程,和Django的Web进程不在同一个内存空间,直接调用它的函数行不通。我们可以用Redis做一个轻量级的消息队列:
- Django在信号触发时,把要发送的MQTT主题、负载等信息放到Redis队列里
monitor.py在运行MQTT客户端的同时,持续监听这个队列,一旦有消息就用已有的客户端实例发送
步骤1:准备依赖
先安装Redis的Python客户端:
pip install redis
确保你的机器上已经安装并运行了Redis服务(大部分Linux发行版可以用包管理器安装,比如apt install redis-server)。
步骤2:修改monitor.py,添加队列监听
更新你的常驻监控客户端,让它同时监听Redis队列:
import paho.mqtt.client as mqtt import redis import json # 初始化Redis客户端 redis_client = redis.Redis(host='localhost', port=6379, db=0) # 初始化你的MQTT客户端(保留你已有的加密、认证配置) mqtt_client = mqtt.Client(client_id="monitor_client") mqtt_client.username_pw_set("abc", "abc") # 配置TLS(用你现有的证书路径) mqtt_client.tls_set(ca_certs="", certfile="", keyfile="") mqtt_client.connect("localhost", 8081) mqtt_client.loop_start() # 启动MQTT后台循环 def handle_mqtt_publish(): """监听Redis队列,用已有的MQTT客户端发送消息""" while True: # 阻塞式获取队列消息,避免空轮询浪费资源 _, raw_msg = redis_client.blpop("mqtt_publish_queue") try: msg_data = json.loads(raw_msg) # 用已连接的MQTT客户端发送消息 result = mqtt_client.publish( topic=msg_data["topic"], payload=msg_data["payload"], retain=msg_data.get("retain", False) ) # 可以添加发送结果的日志记录 result.wait_for_publish() except Exception as e: # 处理异常,比如日志记录或重试 print(f"Failed to send MQTT message: {str(e)}") if __name__ == "__main__": handle_mqtt_publish()
步骤3:修改Django信号处理函数
把原来直接调用publish.single的逻辑改成往Redis队列塞消息:
from django.db.models.signals import post_save from django.dispatch import receiver from .models import Post import redis import json # 初始化Redis客户端 redis_client = redis.Redis(host='localhost', port=6379, db=0) @receiver(post_save, sender=Post) def save_post(sender, instance, **kwargs): # 构造要发送的MQTT负载 message_payload = { "client_id": "abc", "message": f"Created new model: {str(instance)}", } # 封装MQTT发送所需的全部信息 mqtt_task = { "topic": "house/StateServer/receive", "payload": json.dumps(message_payload), "retain": False } # 把任务放到Redis队列 redis_client.rpush("mqtt_publish_queue", json.dumps(mqtt_task))
为什么这个方案更好?
- 复用连接:全程只用你已经配置好的MQTT客户端实例,避免频繁创建销毁连接的开销
- 解耦逻辑:Django只负责生产消息,MQTT发送逻辑由专门的monitor进程处理,符合单一职责原则
- 可靠性:Redis队列会暂存消息,即使monitor临时重启,也不会丢失待发送的消息(如果需要持久化,可以开启Redis的RDB/AOF持久化)
额外注意事项
- 确保
monitor.py是用进程管理器(比如systemd、supervisor)托管的常驻进程,避免意外退出 - 可以根据需求添加日志记录,方便排查问题
- 如果你的部署环境是多台机器,只需要把Redis地址改成公共可访问的即可,无需修改其他逻辑
内容的提问来源于stack exchange,提问作者Ondřej Holík
相关产品推荐
相关产品推荐

