如何在Django+mqttasgi异步环境中实现数据库更新?
解决Django异步MQTT消费者中更新数据库的问题
你的报错是因为Django ORM操作是同步阻塞的,不能直接在异步函数(async def)中调用。必须用Django提供的sync_to_async工具,把同步的数据库操作包装成异步兼容的形式。
修改步骤及代码:
- 导入
sync_to_async工具 - 将数据库查询、更新逻辑封装成独立的同步函数,并用
@sync_to_async装饰 - 在异步的
receive方法中用await调用这个包装后的函数 - 增加异常处理,避免因设备不存在导致服务崩溃
修改后的完整代码:
import json from datetime import datetime from mqttasgi.consumers import MqttConsumer from asgiref.sync import sync_to_async # 新增导入 import django import os os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'myproj.settings') django.setup() from .models import Device # 封装同步数据库操作,用sync_to_async装饰 @sync_to_async def update_device_online(devicename, is_online): try: device = Device.objects.get(name=devicename) device.online = is_online device.save() return True except Device.DoesNotExist: print(f"Warning: 设备 {devicename} 不存在") return False class MyMqttConsumer(MqttConsumer): async def connect(self): await self.subscribe('test/#', 2) async def receive(self, mqtt_message): print(f"{datetime.now()} -- Received a message at topic: { mqtt_message['topic'] }") print(f"With payload: { mqtt_message['payload'] }") print(f"And QOS: { mqtt_message['qos']}") # 增加topic格式校验,避免分割出错 topic_parts = mqtt_message['topic'].split('/', 2) if len(topic_parts) != 3: print(f"无效的topic格式: {mqtt_message['topic']}") return site, devicename, sub_topic = topic_parts if sub_topic == 'online': # 解析payload,转换为布尔值 payload_str = mqtt_message['payload'].decode('UTF-8').strip() is_online = payload_str.lower() != 'false' # 调用异步包装的数据库更新函数 await update_device_online(devicename, is_online) async def disconnect(self): await self.unsubscribe('test/#')
关键说明:
sync_to_async会把同步函数放到单独的线程中执行,避免阻塞异步事件循环- 必须用
await调用被sync_to_async装饰的函数,因为它返回的是协程对象 - 增加异常处理和格式校验,提升消费者的稳定性
内容的提问来源于stack exchange,提问作者Steve Lalonde
相关产品推荐
相关产品推荐

