# Water 项目架构图 本文档按当前项目代码结构整理,重点包含后端模块、智能灌溉业务、MQTT 设备通信、ACK 应答与重发链路。 ## 整体模块 ```mermaid flowchart TB APP["移动端 / 管理端 App"] DEVICE["灌溉设备"] BROKER["MQTT Broker
TLS: ssl://service.reinkun.com:8883"] subgraph BACKEND["water 后端服务"] ADMIN["water-admin
Spring Boot 启动入口"] APPMOD["water-modules/water-app
设备、排程、浇水记录、MQTT业务分发"] SYSTEM["water-modules/water-system
用户、角色、租户、权限"] COMMON["water-common
通用能力"] MQTT["water-common-mqtt
MQTT连接、TLS、订阅、发布、ACK重发"] REDISCOMMON["water-common-redis
Redis 工具与 Redisson"] DBACCESS["water-common-mybatis
MyBatis Plus / 数据权限"] end MYSQL[("MySQL
业务数据")] REDIS[("Redis
缓存、设备状态、待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 请求
/app/v1/**"] CONTROLLER["AppController
设备、排程、日志、用户接口"] SERVICE["Service 层
IAppDeviceService / IAppScheduleService / IAppWateringLogService"] MAPPER["Mapper 层
MyBatis Plus Mapper"] DB[("MySQL")] HTTP --> CONTROLLER CONTROLLER --> SERVICE SERVICE --> MAPPER MAPPER --> DB CONTROLLER --> CMDPUB["IDeviceCommandPublisher
设备命令发布接口"] CMDPUB --> MQTTIMPL["DeviceMqttCommandPublisher
MQTT命令实现"] ``` ## MQTT 高并发消费链路 ```mermaid flowchart TB DEVICE1["设备 A"] DEVICE2["设备 B"] DEVICEN["设备 N"] BROKER["MQTT Broker"] CLIENT["MqttAsyncClient
单个后端订阅客户端"] QUEUE["有界内存队列
queue-capacity: 20000"] CONSUMERS["mqtt-consumer-*
consumer-count: 16
batch-size: 100"] DISPATCHER["MqttMessageDispatcher
topic 路由"] DATA["DeviceDataHandler
/water/{deviceNo}/data"] STATUS["DeviceStatusHandler
/water/{deviceNo}/status"] ERROR["ErromesHandler
/water/{deviceNo}/erromes"] ACK["IDeviceCommandAckHandler
/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}
设备最新在线状态
TTL: 300s"] PENDING["mqtt:command:pending:{commandId}
待ACK命令详情
TTL: 86400s"] PENDING_IDS["mqtt:command:pending:ids
待ACK commandId 集合"] ACK["mqtt:command:ack:{commandId}
设备ACK结果
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`