200 lines
6.0 KiB
Markdown
200 lines
6.0 KiB
Markdown
# Water 项目架构图
|
||
|
||
本文档按当前项目代码结构整理,重点包含后端模块、智能灌溉业务、MQTT 设备通信、ACK 应答与重发链路。
|
||
|
||
## 整体模块
|
||
|
||
```mermaid
|
||
flowchart TB
|
||
APP["移动端 / 管理端 App"]
|
||
DEVICE["灌溉设备"]
|
||
BROKER["MQTT Broker<br/>TLS: ssl://service.reinkun.com:8883"]
|
||
|
||
subgraph BACKEND["water 后端服务"]
|
||
ADMIN["water-admin<br/>Spring Boot 启动入口"]
|
||
APPMOD["water-modules/water-app<br/>设备、排程、浇水记录、MQTT业务分发"]
|
||
SYSTEM["water-modules/water-system<br/>用户、角色、租户、权限"]
|
||
COMMON["water-common<br/>通用能力"]
|
||
MQTT["water-common-mqtt<br/>MQTT连接、TLS、订阅、发布、ACK重发"]
|
||
REDISCOMMON["water-common-redis<br/>Redis 工具与 Redisson"]
|
||
DBACCESS["water-common-mybatis<br/>MyBatis Plus / 数据权限"]
|
||
end
|
||
|
||
MYSQL[("MySQL<br/>业务数据")]
|
||
REDIS[("Redis<br/>缓存、设备状态、待ACK命令")]
|
||
|
||
APP -->|"HTTP REST"| ADMIN
|
||
ADMIN --> APPMOD
|
||
ADMIN --> SYSTEM
|
||
ADMIN --> COMMON
|
||
APPMOD --> DBACCESS
|
||
SYSTEM --> DBACCESS
|
||
DBACCESS --> MYSQL
|
||
APPMOD --> REDISCOMMON
|
||
MQTT --> REDISCOMMON
|
||
REDISCOMMON --> REDIS
|
||
MQTT <-->|"MQTT over TLS"| BROKER
|
||
DEVICE <-->|"MQTT over TLS"| BROKER
|
||
MQTT --> APPMOD
|
||
```
|
||
|
||
## 后端分层
|
||
|
||
```mermaid
|
||
flowchart LR
|
||
HTTP["HTTP 请求<br/>/app/v1/**"]
|
||
CONTROLLER["AppController<br/>设备、排程、日志、用户接口"]
|
||
SERVICE["Service 层<br/>IAppDeviceService / IAppScheduleService / IAppWateringLogService"]
|
||
MAPPER["Mapper 层<br/>MyBatis Plus Mapper"]
|
||
DB[("MySQL")]
|
||
|
||
HTTP --> CONTROLLER
|
||
CONTROLLER --> SERVICE
|
||
SERVICE --> MAPPER
|
||
MAPPER --> DB
|
||
|
||
CONTROLLER --> CMDPUB["IDeviceCommandPublisher<br/>设备命令发布接口"]
|
||
CMDPUB --> MQTTIMPL["DeviceMqttCommandPublisher<br/>MQTT命令实现"]
|
||
```
|
||
|
||
## MQTT 高并发消费链路
|
||
|
||
```mermaid
|
||
flowchart TB
|
||
DEVICE1["设备 A"]
|
||
DEVICE2["设备 B"]
|
||
DEVICEN["设备 N"]
|
||
BROKER["MQTT Broker"]
|
||
CLIENT["MqttAsyncClient<br/>单个后端订阅客户端"]
|
||
QUEUE["有界内存队列<br/>queue-capacity: 20000"]
|
||
CONSUMERS["mqtt-consumer-*<br/>consumer-count: 16<br/>batch-size: 100"]
|
||
DISPATCHER["MqttMessageDispatcher<br/>topic 路由"]
|
||
|
||
DATA["DeviceDataHandler<br/>/water/{deviceNo}/data"]
|
||
STATUS["DeviceStatusHandler<br/>/water/{deviceNo}/status"]
|
||
ERROR["ErromesHandler<br/>/water/{deviceNo}/erromes"]
|
||
ACK["IDeviceCommandAckHandler<br/>/water/{deviceNo}/ack"]
|
||
|
||
DEVICE1 --> BROKER
|
||
DEVICE2 --> BROKER
|
||
DEVICEN --> BROKER
|
||
BROKER --> CLIENT
|
||
CLIENT -->|"快速入队"| QUEUE
|
||
QUEUE -->|"批量拉取"| CONSUMERS
|
||
CONSUMERS --> DISPATCHER
|
||
DISPATCHER --> DATA
|
||
DISPATCHER --> STATUS
|
||
DISPATCHER --> ERROR
|
||
DISPATCHER --> ACK
|
||
```
|
||
|
||
## 设备命令 ACK 与重发
|
||
|
||
```mermaid
|
||
sequenceDiagram
|
||
participant App as 移动端 App
|
||
participant API as AppController
|
||
participant Pub as DeviceMqttCommandPublisher
|
||
participant Redis as Redis
|
||
participant MQTT as MqttClientManager
|
||
participant Broker as MQTT Broker
|
||
participant Dev as 设备
|
||
participant Ack as MqttCommandAckService
|
||
participant Retry as MqttCommandRetryTask
|
||
|
||
App->>API: switchDevice(workStatus, durationMin)
|
||
API->>Pub: send(DeviceCommand)
|
||
Pub->>Redis: 保存 pending commandId
|
||
Pub->>MQTT: publish /water/{deviceNo}/command
|
||
MQTT->>Broker: 下发命令
|
||
Broker->>Dev: 命令到达设备
|
||
|
||
alt 设备正常回复
|
||
Dev->>Broker: publish /water/{deviceNo}/ack
|
||
Broker->>MQTT: ACK 消息
|
||
MQTT->>Ack: handleAck(deviceNo, payload)
|
||
Ack->>Redis: 删除 pending,保存 ack 结果
|
||
else 设备断网或未回复
|
||
Retry->>Redis: 扫描 pending 命令
|
||
Retry->>Redis: 检查设备状态缓存
|
||
alt 设备在线且未超过重试次数
|
||
Retry->>MQTT: 重新 publish 命令
|
||
MQTT->>Broker: 重新下发
|
||
else 设备离线
|
||
Retry->>Redis: 延后 nextRetryAt
|
||
else 超过最大重试次数
|
||
Retry->>Redis: 删除 pending
|
||
end
|
||
end
|
||
```
|
||
|
||
## Redis Key 规划
|
||
|
||
```mermaid
|
||
flowchart TB
|
||
REDIS[("Redis")]
|
||
STATUS["mqtt:device:status:{deviceNo}<br/>设备电量在线心跳<br/>TTL: 600s"]
|
||
PENDING["mqtt:command:pending:{commandId}<br/>待ACK命令详情<br/>TTL: 86400s"]
|
||
PENDING_IDS["mqtt:command:pending:ids<br/>待ACK commandId 集合"]
|
||
ACK["mqtt:command:ack:{commandId}<br/>设备ACK结果<br/>TTL: 86400s"]
|
||
|
||
REDIS --> STATUS
|
||
REDIS --> PENDING
|
||
REDIS --> PENDING_IDS
|
||
REDIS --> ACK
|
||
```
|
||
|
||
## 核心 Topic 约定
|
||
|
||
```mermaid
|
||
flowchart LR
|
||
DEVICE["设备"]
|
||
SERVER["后端"]
|
||
|
||
DEVICE -->|"/water/{deviceNo}/data"| SERVER
|
||
DEVICE -->|"/water/{deviceNo}/status"| SERVER
|
||
DEVICE -->|"/water/{deviceNo}/erromes"| SERVER
|
||
DEVICE -->|"/water/{deviceNo}/ack"| SERVER
|
||
SERVER -->|"/water/{deviceNo}/command"| DEVICE
|
||
```
|
||
|
||
设备 ACK 示例:
|
||
|
||
```json
|
||
{
|
||
"commandId": "后端下发的commandId",
|
||
"status": "success",
|
||
"message": "ok"
|
||
}
|
||
```
|
||
|
||
## 高并发关注点
|
||
|
||
```mermaid
|
||
flowchart TB
|
||
A["上千设备并发连接"] --> B["MQTT Broker 承载连接数和 TLS 握手"]
|
||
B --> C["后端单/少量订阅客户端消费通配 topic"]
|
||
C --> D["有界队列吸收突发流量"]
|
||
D --> E["批量消费降低线程调度成本"]
|
||
E --> F["电量心跳写 Redis 并同步设备状态"]
|
||
F --> G["数据/告警后续建议批量落库"]
|
||
```
|
||
|
||
当前已完成的关键优化:
|
||
|
||
- MQTT 使用 `MqttAsyncClient`
|
||
- MQTT 连接支持 TLS
|
||
- MQTT 回调只入队,不直接执行业务
|
||
- 消息按批次消费
|
||
- 设备状态写 Redis 缓存
|
||
- 设备命令支持 ACK
|
||
- 未 ACK 命令支持离线延后和超时重发
|
||
- 日志异步队列已扩大
|
||
|
||
后续建议:
|
||
|
||
- 设备数据和告警数据落库改为批量写入
|
||
- 对设备命令增加业务状态表,便于前端查询命令执行状态
|
||
- MQTT Broker 使用集群或至少做连接数、会话数、消息速率监控
|
||
- 生产环境 TLS 不建议开启 `skip-verify`
|