提交 22eb5e22 authored 作者: 陈泽健's avatar 陈泽健

docs(device-sim): 消息记录架构重构 PRD + 执行计划 + HANDOFF 更新

- 新增 PRD: 消息记录架构重构与实时消息流

- 新增执行计划: 消息队列实现方案

- 更新 HANDOFF: 本次会话进度和已知问题
Co-Authored-By: 's avatarClaude <noreply@anthropic.com>
上级 18fbe299
# PRD 需求优化 — 消息记录架构重构与实时消息流
> **文档类型**: PRD 需求优化文档
> **创建日期**: 2026-08-05
> **作者**: czj
> **优先级**: P0
> **版本**: v1.0
> **状态**: 待评审
> **关联文档**: `_PRD_需求文档_设备模拟模块.md`、`_PRD_需求优化_环境配置多主题与设备接收记录.md`
---
## 一、需求背景
### 1.1 当前问题
设备模拟模块存在以下架构问题:
| # | 问题 | 根因 | 影响 |
|---|------|------|------|
| 1 | 上报记录数据库写入失败 | `_sync_report_log()` 在模拟器线程中运行,数据库 session 不兼容线程 | `totalReports` 始终为 0,上报记录为空 |
| 2 | WebSocket 实时推送失败 | 模拟器线程中没有正确的事件循环,回调无法执行 | 消息流页面显示"暂无消息" |
| 3 | 消息记录与推送耦合 | 数据库写入和 WebSocket 推送在同一个异步调用链中 | 一个失败导致另一个也失败 |
**核心矛盾**:模拟器在独立线程中运行(`sync_playwright` + `threading`),而数据库操作和 WebSocket 推送需要异步事件循环。当前实现通过 `asyncio.new_event_loop()` 在线程中创建新循环,但数据库 session 不支持跨线程使用,导致 `Session's transaction has been rolled back` 错误。
### 1.2 目标用户
| 用户角色 | 使用场景 | 核心需求 |
|---------|---------|---------|
| 测试工程师 | 模拟设备上报/监听 | 实时查看上报和下发消息 |
| 开发人员 | 调试 MQTT 通信 | 消息流实时显示,支持搜索过滤 |
| 运维人员 | 批量管理设备 | 批量操作状态实时反馈 |
### 1.3 需求目标
1. **修复上报记录数据库写入** — 模拟器上报消息后能正确写入数据库
2. **实现 WebSocket 实时推送** — 消息流页面实时显示设备上报/下发消息
3. **解耦消息记录与推送** — 架构层面隔离,互不影响
4. **保证 MQTT 上报不受影响** — 消息记录和推送失败不影响 MQTT 上报
---
## 二、功能需求
### 2.1 消息队列中间层
**功能描述**:在模拟器线程和数据库/推送之间引入消息队列,解耦异步操作。
**功能清单**
| 功能 | 说明 | 优先级 |
|------|------|--------|
| 消息入队 | 模拟器上报/接收消息时,将消息写入内存队列 | P0 |
| 消息消费 | 主线程异步消费队列,写入数据库 | P0 |
| WebSocket 推送 | 消费消息时同时推送到前端 | P0 |
| 队列溢出保护 | 队列长度超限时丢弃旧消息 | P1 |
| 消息统计 | 记录入队/出队/丢弃数量 | P2 |
**消息队列数据结构**
```python
@dataclass
class QueueMessage:
device_id: str # 设备 ID(数据库主键)
device_name: str # 设备名称
device_type: str # 设备类型
topic: str # MQTT 主题
payload: dict # 消息内容
direction: str # publish/subscribe
status: str # success/failed
error_message: str # 错误信息
timestamp: datetime # 消息时间
```
### 2.2 实时消息流修复
**功能描述**:修复 WebSocket 推送,使消息流页面能实时显示设备消息。
**功能清单**
| 功能 | 说明 | 优先级 |
|------|------|--------|
| WebSocket 连接 | 前端连接 `ws://host/api/device-sim/ws/messages` | P0 |
| 消息推送 | 消费队列时推送消息到所有连接的客户端 | P0 |
| 按设备类型过滤 | 推送时包含 device_type,前端按类型过滤 | P0 |
| 连接管理 | 自动清理断开的 WebSocket 连接 | P1 |
| 心跳保活 | 支持 ping/pong 心跳 | P2 |
---
## 三、非功能需求
### 3.1 性能要求
| 指标 | 要求 | 说明 |
|------|------|------|
| 队列消费延迟 | < 100ms | 消息从入队到推送完成 |
| 队列容量 | 最多 10000 条 | 超出后丢弃最旧消息 |
| WebSocket 并发 | 支持 10+ 客户端 | 多人同时查看消息流 |
| 对 MQTT 上报影响 | 0 | 消息记录失败不影响上报 |
### 3.2 可靠性要求
- 消息队列是内存队列,程序重启时未消费的消息会丢失(可接受)
- MQTT 上报不受消息队列和数据库影响
- WebSocket 推送失败不影响数据库写入
---
## 四、技术方案
### 4.1 架构设计
```
┌─────────────────────────────────────────────────────────────┐
│ 模拟器线程 │
│ ┌──────────┐ ┌──────────┐ ┌──────────────────────┐ │
│ │ MQTT 上报 │───▶│ 回调触发 │───▶│ 消息入队(非阻塞) │ │
│ └──────────┘ └──────────┘ └──────────┬───────────┘ │
│ │ │
└─────────────────────────────────────────────┼──────────────┘
┌─────────────────────────────────────────────────────────────┐
│ 主线程异步循环 │
│ ┌──────────────────┐ ┌──────────┐ ┌──────────────┐ │
│ │ 消息出队(异步) │───▶│ 数据库写入│ │ │ │
│ └──────────────────┘ └──────────┘ │ WebSocket 推送│ │
│ ┌──────────┐ │ │ │
│ │ 更新统计 │───▶│ │ │
│ └──────────┘ └──────────────┘ │
└─────────────────────────────────────────────────────────────┘
```
### 4.2 核心实现
**消息队列**:使用 `asyncio.Queue`,在主线程异步循环中消费。
**关键改动**
- 模拟器回调只做入队操作(`queue.put_nowait()`),不涉及数据库或异步调用
- 消费循环在 `main.py``lifespan` 中启动,使用 `asyncio.create_task()`
- 消费循环中处理数据库写入和 WebSocket 推送
### 4.3 改动范围
| 文件 | 变更类型 | 说明 |
|------|----------|------|
| `backend/app/services/device_sim_service.py` | 修改 | 重构消息记录机制,使用队列替代直接写入 |
| `backend/app/routers/device_sim.py` | 修改 | WebSocket 推送逻辑调整 |
| `backend/app/main.py` | 修改 | 启动消息消费循环 |
---
## 五、约束与限制
- **不改动模拟器线程模型** — 模拟器必须使用同步 API(Windows 限制)
- **不改动数据库模型** — 表结构不变,只改写入方式
- **不改动前端代码** — 消息流组件已实现,只需后端推送正常
---
## 六、验收标准
| 验收项 | 预期结果 |
|--------|---------|
| 启动设备后,消息流页面实时显示上报消息 | 上报消息每 2 秒刷新一次 |
| `totalReports` 字段正确递增 | 每次上报后 +1 |
| 上报记录列表可查询 | 按设备/方向/时间筛选正常 |
| 消息流显示下发消息 | 模拟器收到指令时显示 subscribe 方向消息 |
| 停止设备后消息流停止 | 无新消息推入 |
| 批量操作后消息流反映状态变化 | 启动/停止设备后消息流响应 |
---
## 七、风险与应对
| 风险 | 影响 | 概率 | 应对 |
|------|------|------|------|
| 消费循环异常导致消息积压 | 内存占用增加 | 低 | 加入异常捕获和队列长度限制 |
| 程序重启时队列消息丢失 | 少量上报记录缺失 | 中 | 可接受,MQTT 上报不受影响 |
| WebSocket 推送阻塞消费循环 | 数据库写入延迟 | 低 | 推送使用 `create_task` 不阻塞 |
---
*本文档为消息记录架构重构与实时消息流的需求文档,供执行计划参考。*
# 执行计划 — 消息记录架构重构与实时消息流
> **文档类型**: 执行计划文档
> **创建日期**: 2026-08-05
> **关联 PRD**: `_PRD_需求优化_消息记录架构重构与实时消息流.md` v1.0
> **负责人**: czj
> **优先级**: P0
> **状态**: 待评审
---
## 一、执行概述
### 1.1 需求概要
重构设备模拟模块的消息记录机制,引入消息队列解耦模拟器线程和数据库/推送操作:
- 模拟器线程:只负责 MQTT 上报 + 消息入队(非阻塞)
- 主线程异步循环:消费队列 + 写入数据库 + WebSocket 推送
### 1.2 技术变更范围
| 层级 | 影响文件 | 变更类型 |
|------|----------|----------|
| 服务层 | `backend/app/services/device_sim_service.py` | 重构消息记录机制 |
| 路由层 | `backend/app/routers/device_sim.py` | WebSocket 推送调整 |
| 入口层 | `backend/app/main.py` | 启动消费循环 |
### 1.3 预计工期
| 阶段 | 内容 | 预计工时 |
|------|------|---------|
| 后端重构 | 消息队列 + 消费循环 | 2 小时 |
| 测试验证 | 功能测试 + 边界测试 | 1 小时 |
| 文档更新 | HANDOFF + 部署 | 0.5 小时 |
---
## 二、任务分解
### 2.1 消息队列实现
#### 任务清单
| 任务ID | 任务名称 | 文件 | 说明 | 优先级 |
|--------|----------|------|------|--------|
| T1.1 | 定义消息数据结构 | `device_sim_service.py` | `QueueMessage` dataclass | P0 |
| T1.2 | 创建全局消息队列 | `device_sim_service.py` | `asyncio.Queue(maxsize=10000)` | P0 |
| T1.3 | 修改模拟器回调 | `device_sim_service.py` | 入队替代直接写入 | P0 |
#### 实施代码
```python
# device_sim_service.py
from dataclasses import dataclass
from datetime import datetime
import asyncio
@dataclass
class QueueMessage:
device_id: str
device_name: str
device_type: str
topic: str
payload: dict
direction: str
status: str
error_message: str
timestamp: datetime
# 全局消息队列
_message_queue: asyncio.Queue = asyncio.Queue(maxsize=10000)
def _report_callback(self, device_id: str, topic: str, payload: dict,
direction: str = "publish", status: str = "success",
error_message: str = None):
"""模拟器回调(在模拟器线程中执行)— 只入队,不阻塞"""
msg = QueueMessage(
device_id=device_id,
device_name=self.device_name,
device_type=self.device_type,
topic=topic,
payload=payload,
direction=direction,
status=status,
error_message=error_message,
timestamp=datetime.now(),
)
try:
_message_queue.put_nowait(msg)
except asyncio.QueueFull:
logger.warning(f"消息队列已满,丢弃消息: {device_id}")
```
---
### 2.2 消费循环实现
#### 任务清单
| 任务ID | 任务名称 | 文件 | 说明 | 优先级 |
|--------|----------|------|------|--------|
| T2.1 | 实现消费循环 | `device_sim_service.py` | 从队列取消息写入数据库 | P0 |
| T2.2 | WebSocket 推送集成 | `device_sim_service.py` | 消费时触发推送 | P0 |
| T2.3 | 启动消费循环 | `main.py` | lifespan 中 create_task | P0 |
#### 实施代码
```python
# device_sim_service.py
async def _message_consumer():
"""消息消费循环(在主线程异步循环中运行)"""
while True:
try:
msg: QueueMessage = await _message_queue.get()
# 写入数据库
log = ReportLog(
id=generate_id("log"),
device_id=msg.device_id,
topic=msg.topic,
payload=msg.payload,
direction=msg.direction,
status=msg.status,
error_message=msg.error_message,
)
# 需要获取数据库 session
# WebSocket 推送
await _broadcast_device_message({
"event": "device_message",
"data": {...}
})
except Exception as e:
logger.error(f"消息消费异常: {e}")
await asyncio.sleep(1)
```
```python
# main.py
@asynccontextmanager
async def lifespan(app: FastAPI):
# ... 现有初始化代码 ...
# 启动消息消费循环
from app.services.device_sim_service import _message_consumer
consumer_task = asyncio.create_task(_message_consumer())
yield
# 清理
consumer_task.cancel()
```
---
### 2.3 WebSocket 推送修复
#### 任务清单
| 任务ID | 任务名称 | 文件 | 说明 | 优先级 |
|--------|----------|------|------|--------|
| T3.1 | 调整推送逻辑 | `device_sim.py` | 移除线程不安全的回调 | P0 |
| T3.2 | 保留连接管理 | `device_sim.py` | WebSocket 端点和广播函数 | P0 |
---
## 三、技术要点
### 3.1 数据库 Session 问题
**问题**:消费循环需要数据库 session,但 `asyncio.Queue` 的消费者在不同上下文中。
**方案**:消费循环中使用 `async for` 从队列获取消息,然后通过 `get_db()` 获取新 session:
```python
async def _message_consumer():
async with get_db_session() as db:
while True:
msg = await _message_queue.get()
service = DeviceSimService(db)
await service._create_report_log_from_queue(msg)
```
### 3.2 异常处理
消费循环中的异常不应中断循环,需捕获后继续:
```python
except Exception as e:
logger.error(f"消息消费异常: {e}")
await asyncio.sleep(1) # 避免快速失败循环
continue
```
---
## 四、部署方案
### 4.1 本地开发
```bash
cd backend
uvicorn app.main:app --reload --port 8001
```
### 4.2 服务器部署
```bash
# 本地构建
cd frontend && npm run build
# 上传到服务器
python deploy_script.py
# 重启容器
docker restart plat-auto-test-app
```
---
## 五、验收检查清单
| 验收项 | 检查方式 | 预期结果 |
|--------|---------|---------|
| 启动设备,查看消息流 | 打开门口屏模拟页面 | 消息实时显示 |
| 检查上报记录列表 | 刷新记录列表 | 显示新记录 |
| 检查 `totalReports` | 查看设备详情 | 数值正确 |
| 停止设备 | 点击停止按钮 | 消息流停止更新 |
| 查看服务器日志 | `docker logs` | 无 `Session's transaction` 错误 |
---
## 六、风险与应对
| 风险 | 影响 | 概率 | 应对方案 |
|------|------|------|---------|
| 消费循环异常 | 消息积压 | 低 | 捕获异常 + 队列长度限制 |
| 性能下降 | 延迟增加 | 低 | 监控队列长度,优化消费速度 |
---
*本文档为消息记录架构重构与实时消息流的执行计划,需在评审后实施。*
\ No newline at end of file
此差异已折叠。
Markdown 格式
0%
您添加了 0 到此讨论。请谨慎行事。
请先完成此评论的编辑!
注册 或者 后发表评论