Django应用整合远程噪音计查询代码,实现多设备并发实时更新前端数据
Django整合多台噪音计实时查询与WebSockets推送方案建议
核心思路
要实现多设备并发查询+实时前端推送,需解决两个核心问题:异步并发的设备管理和与Django Channels的数据流打通。关键是将单设备查询逻辑重构为可复用的异步类,用asyncio原生调度实现多设备并发,再通过消息队列将数据传递给Channels消费者推送到前端。
1. 重构噪音计查询代码为可复用异步类
将原单设备函数改为类,每个实例对应一台噪音计,管理连接状态、数据读取,并替换阻塞的time.sleep为异步的asyncio.sleep,避免阻塞事件循环影响并发。
import datetime import asyncio from typing import Optional class NoiseMeterClient: def __init__(self, host: str, port: int, meter_id: str): self.host = host self.port = port self.meter_id = meter_id # 设备唯一标识,用于前端区分 self.reader: Optional[asyncio.StreamReader] = None self.writer: Optional[asyncio.StreamWriter] = None self.running = False self.message_queue = asyncio.Queue() # 存储待推送的数据 async def connect(self): """建立连接并初始化设备""" try: self.reader, self.writer = await asyncio.open_connection(self.host, self.port) # 设备重置 await self.send_command('*RST\n') await asyncio.sleep(1) # 查询设备ID resp = await self.send_command('*IDN?\n') print(f"[{self.meter_id}] 设备ID: {resp!r}") # 开启自动保存 resp = await self.send_command('SYST:KEY NEXT NEXT NEXT ENTER PREV ENTER ESC\n') print(f"[{self.meter_id}] 自动保存状态: {resp!r}") # 启动测量 await self.send_command('INIT START\n') await asyncio.sleep(6) self.running = True print(f"[{self.meter_id}] 初始化完成,开始读取数据") except Exception as e: print(f"[{self.meter_id}] 连接/初始化失败: {str(e)}") self.running = False if self.writer: self.writer.close() await self.writer.wait_closed() async def send_command(self, command: str) -> str: """发送命令并返回响应""" if not self.writer: raise ConnectionError("未建立设备连接") self.writer.write(command.encode()) await self.writer.drain() data = await self.reader.read(1024) return data.decode().strip() async def start_reading(self, interval: float = 1.0): """按指定间隔读取数据,存入消息队列""" if not self.running: await self.connect() count = 0 try: while self.running: # 触发测量+查询数据 await self.send_command('MEAS:INIT\n') laeq_value = await self.send_command('MEAS:SLM:123? LAEQ\n') # 构造结构化数据 data = { 'meter_id': self.meter_id, 'laeq': laeq_value, 'timestamp': datetime.datetime.now().isoformat() } await self.message_queue.put(data) print(f"[{self.meter_id}] 读取数据: {data}") await asyncio.sleep(interval) count += 1 if count > 15: # 可根据需求调整停止条件,或移除该逻辑 await self.stop() except Exception as e: print(f"[{self.meter_id}] 读取失败: {str(e)}") await self.stop() async def stop(self): """停止测量并关闭连接""" self.running = False if self.writer: await self.send_command('INIT STOP\n') self.writer.close() await self.writer.wait_closed() print(f"[{self.meter_id}] 连接已关闭")
2. 整合Django Channels实现数据推送
通过AsyncWebsocketConsumer监听设备消息队列,将数据推送到前端。用全局字典管理多设备实例,确保同一设备的多个前端连接共享一个查询任务。
消费者代码(consumers.py)
import asyncio import json from channels.generic.websocket import AsyncWebsocketConsumer from .noise_meter import NoiseMeterClient # 全局设备实例管理器,避免重复创建连接 meter_clients = {} class NoiseMeterConsumer(AsyncWebsocketConsumer): async def connect(self): # 从URL参数获取设备ID self.meter_id = self.scope['url_route']['kwargs']['meter_id'] self.room_group_name = f'noise_meter_{self.meter_id}' # 加入设备对应的推送组 await self.channel_layer.group_add( self.room_group_name, self.channel_name ) await self.accept() # 初始化设备客户端(未存在时) if self.meter_id not in meter_clients: # 建议从数据库/配置文件读取设备的host和port,此处为示例硬编码 host = '192.168.10.1' port = 6889 client = NoiseMeterClient(host, port, self.meter_id) meter_clients[self.meter_id] = client # 启动数据读取+推送任务 asyncio.create_task(self._forward_data(client)) async def disconnect(self, close_code): # 离开推送组 await self.channel_layer.group_discard( self.room_group_name, self.channel_name ) # 清理设备连接(无活跃前端时) if self.meter_id in meter_clients: client = meter_clients[self.meter_id] await client.stop() del meter_clients[self.meter_id] async def _forward_data(self, client: NoiseMeterClient): """从设备队列取数据,推送到前端组""" await client.start_reading() while client.running: data = await client.message_queue.get() await self.channel_layer.group_send( self.room_group_name, { 'type': 'noise_data_message', 'data': data } ) async def noise_data_message(self, event): """将数据发送到Websocket""" data = event['data'] await self.send(text_data=json.dumps(data))
Channels路由配置(routing.py)
from django.urls import re_path from . import consumers websocket_urlpatterns = [ re_path(r'ws/noise-meter/(?P<meter_id>\w+)/$', consumers.NoiseMeterConsumer.as_asgi()), ]
3. 前端Websocket连接示例
// 可从页面参数/用户选择获取设备ID const meterId = 'office_meter'; const ws = new WebSocket(`ws://${window.location.host}/ws/noise-meter/${meterId}/`); ws.onmessage = function(event) { const data = JSON.parse(event.data); // 更新页面DOM document.getElementById('meter-id').textContent = data.meter_id; document.getElementById('laeq-value').textContent = `${data.laeq} dB`; document.getElementById('timestamp').textContent = data.timestamp; }; ws.onerror = function(error) { console.error('WebSocket连接错误:', error); }; ws.onclose = function() { console.log('WebSocket连接已关闭'); };
4. 多设备并发支持说明
- 每个设备对应独立的
NoiseMeterClient实例,由meter_clients字典统一管理,不同前端连接可订阅同一设备的数据流。 - 所有操作基于asyncio异步调度,多设备的查询任务会被事件循环自动分配执行,不会互相阻塞。
5. 关键注意事项
- 替换阻塞调用:必须用
asyncio.sleep替代原代码中的time.sleep,否则会阻塞整个异步事件循环,导致多设备并发失效。 - 配置解耦:不要硬编码设备的host和port,建议存入Django数据库或配置文件,前端选择设备时从后端动态获取。
- 异常重连:可在
start_reading方法中加入重连逻辑,当设备断开时自动尝试重新连接。 - 资源清理:确保Websocket断开时停止对应设备的查询任务并关闭连接,避免资源泄漏。
内容的提问来源于stack exchange,提问作者Voxie
相关产品推荐
相关产品推荐

