- 实现REST API服务器,支持IoT数据、手动录入和批量导入接口 - 实现WebSocket服务器,支持实时数据推送和连接管理 - 实现批量导入模块,支持CSV、Excel、JSON多种格式 - 实现数据处理工具,包含字段映射和单位转换功能 - 实现数据模型定义和数据验证机制 - 创建主程序入口和配置文件 - 添加详细的使用文档和API说明
214 lines
7.3 KiB
Python
214 lines
7.3 KiB
Python
"""
|
|
WebSocket 实时数据推送服务器
|
|
支持实时数据推送、连接管理和数据广播
|
|
"""
|
|
import asyncio
|
|
import json
|
|
import websockets
|
|
from datetime import datetime
|
|
from typing import Set, Dict, Any
|
|
import logging
|
|
|
|
# 配置日志
|
|
logging.basicConfig(level=logging.INFO)
|
|
logger = logging.getLogger(__name__)
|
|
|
|
class WebSocketServer:
|
|
"""WebSocket服务器类"""
|
|
|
|
def __init__(self, host: str = "0.0.0.0", port: int = 8765):
|
|
self.host = host
|
|
self.port = port
|
|
self.clients: Set[websockets.WebSocketServerProtocol] = set()
|
|
self.data_history: list = [] # 存储最近的数据用于新连接
|
|
|
|
async def register_client(self, websocket: websockets.WebSocketServerProtocol):
|
|
"""注册新客户端"""
|
|
self.clients.add(websocket)
|
|
client_ip = websocket.remote_address[0]
|
|
logger.info(f"新客户端连接: {client_ip}")
|
|
|
|
# 发送历史数据给新连接的客户端
|
|
if self.data_history:
|
|
await websocket.send(json.dumps({
|
|
"type": "history",
|
|
"data": self.data_history[-50:] # 发送最近50条数据
|
|
}))
|
|
|
|
# 发送欢迎消息
|
|
await websocket.send(json.dumps({
|
|
"type": "welcome",
|
|
"message": "已连接到水务管理系统实时数据服务器",
|
|
"timestamp": datetime.now().isoformat()
|
|
}))
|
|
|
|
async def unregister_client(self, websocket: websockets.WebSocketServerProtocol):
|
|
"""注销客户端"""
|
|
if websocket in self.clients:
|
|
self.clients.remove(websocket)
|
|
client_ip = websocket.remote_address[0]
|
|
logger.info(f"客户端断开连接: {client_ip}")
|
|
|
|
async def broadcast_data(self, data: Dict[str, Any]):
|
|
"""广播数据到所有连接的客户端"""
|
|
if not self.clients:
|
|
return
|
|
|
|
# 添加时间戳
|
|
data["timestamp"] = datetime.now().isoformat()
|
|
|
|
# 保存历史数据
|
|
self.data_history.append(data)
|
|
if len(self.data_history) > 1000: # 只保留最近1000条记录
|
|
self.data_history.pop(0)
|
|
|
|
# 广播数据
|
|
message = json.dumps(data)
|
|
disconnected_clients = []
|
|
|
|
for client in self.clients:
|
|
try:
|
|
await client.send(message)
|
|
except websockets.exceptions.ConnectionClosed:
|
|
disconnected_clients.append(client)
|
|
|
|
# 清理已断开的连接
|
|
for client in disconnected_clients:
|
|
await self.unregister_client(client)
|
|
|
|
async def handle_client_message(self, websocket: websockets.WebSocketServerProtocol, message: str):
|
|
"""处理客户端消息"""
|
|
try:
|
|
data = json.loads(message)
|
|
|
|
if data.get("type") == "subscribe":
|
|
# 处理订阅请求
|
|
subscription_type = data.get("subscription", "all")
|
|
response = {
|
|
"type": "subscription_ack",
|
|
"subscription": subscription_type,
|
|
"message": f"已订阅 {subscription_type} 类型数据"
|
|
}
|
|
await websocket.send(json.dumps(response))
|
|
logger.info(f"客户端订阅了 {subscription_type} 类型数据")
|
|
|
|
elif data.get("type") == "ping":
|
|
# 响应心跳检测
|
|
response = {
|
|
"type": "pong",
|
|
"timestamp": datetime.now().isoformat()
|
|
}
|
|
await websocket.send(json.dumps(response))
|
|
|
|
else:
|
|
logger.warning(f"未知的消息类型: {data.get('type', 'unknown')}")
|
|
|
|
except json.JSONDecodeError:
|
|
logger.error("无效的JSON消息")
|
|
except Exception as e:
|
|
logger.error(f"处理客户端消息时出错: {str(e)}")
|
|
|
|
async def client_handler(self, websocket: websockets.WebSocketServerProtocol, path: str):
|
|
"""处理客户端连接"""
|
|
await self.register_client(websocket)
|
|
|
|
try:
|
|
async for message in websocket:
|
|
await self.handle_client_message(websocket, message)
|
|
except websockets.exceptions.ConnectionClosed:
|
|
pass
|
|
finally:
|
|
await self.unregister_client(websocket)
|
|
|
|
async def start_server(self):
|
|
"""启动WebSocket服务器"""
|
|
logger.info(f"启动WebSocket服务器: {self.host}:{self.port}")
|
|
|
|
# 创建并启动服务器
|
|
self.server = await websockets.serve(
|
|
self.client_handler,
|
|
self.host,
|
|
self.port
|
|
)
|
|
|
|
logger.info("WebSocket服务器已启动")
|
|
return self.server
|
|
|
|
async def send_sensor_data(self, sensor_data: Dict[str, Any]):
|
|
"""发送传感器数据"""
|
|
data = {
|
|
"type": "sensor_data",
|
|
"data_type": sensor_data.get("data_type"),
|
|
"device_id": sensor_data.get("device_id"),
|
|
"value": sensor_data.get("value"),
|
|
"location": sensor_data.get("location"),
|
|
"timestamp": datetime.now().isoformat()
|
|
}
|
|
await self.broadcast_data(data)
|
|
|
|
async def send_alert(self, alert_data: Dict[str, Any]):
|
|
"""发送警报信息"""
|
|
data = {
|
|
"type": "alert",
|
|
"level": alert_data.get("level", "warning"),
|
|
"message": alert_data.get("message"),
|
|
"device_id": alert_data.get("device_id"),
|
|
"timestamp": datetime.now().isoformat()
|
|
}
|
|
await self.broadcast_data(data)
|
|
|
|
# 全局WebSocket服务器实例
|
|
websocket_server = WebSocketServer()
|
|
|
|
# 示例数据生成器
|
|
async def data_generator():
|
|
"""模拟数据生成器"""
|
|
import random
|
|
|
|
while True:
|
|
await asyncio.sleep(5) # 每5秒发送一次数据
|
|
|
|
# 模拟不同的传感器数据
|
|
sensor_types = ["LL", "YL", "SW", "ZD"]
|
|
sensor_type = random.choice(sensor_types)
|
|
|
|
# 根据传感器类型生成合理的数值范围
|
|
if sensor_type == "LL": # 流量
|
|
value = random.uniform(10, 100)
|
|
elif sensor_type == "YL": # 压力
|
|
value = random.uniform(0.1, 1.0)
|
|
elif sensor_type == "SW": # 水位
|
|
value = random.uniform(0, 10)
|
|
else: # ZD 浊度
|
|
value = random.uniform(0, 50)
|
|
|
|
sensor_data = {
|
|
"data_type": sensor_type,
|
|
"device_id": f"device_{random.randint(1, 10)}",
|
|
"value": round(value, 2),
|
|
"location": random.choice(["A区", "B区", "C区", "D区"])
|
|
}
|
|
|
|
await websocket_server.send_sensor_data(sensor_data)
|
|
|
|
# 启动服务器和生成器
|
|
async def main():
|
|
"""主函数"""
|
|
# 启动WebSocket服务器
|
|
server = await websocket_server.start_server()
|
|
|
|
# 启动数据生成器
|
|
generator_task = asyncio.create_task(data_generator())
|
|
|
|
# 保持服务器运行
|
|
try:
|
|
await asyncio.Future() # 永远等待
|
|
except KeyboardInterrupt:
|
|
logger.info("收到中断信号,正在关闭服务器...")
|
|
server.close()
|
|
await server.wait_closed()
|
|
generator_task.cancel()
|
|
await generator_task
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(main()) |