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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 20:47:27