Files
water/docs/project-architecture.md
yuhaiming 70e2356f22 fix(app): 修复 AppController 安全与查询问题
- 加强异常处理、类型安全和图片上传校验
- 优化设备相关查询,避免重复访问数据源
- 补充并记录并发测试与审查修复实施计划
2026-07-17 08:20:44 +08:00

6.0 KiB
Raw Blame History

Water 项目架构图

本文档按当前项目代码结构整理重点包含后端模块、智能灌溉业务、MQTT 设备通信、ACK 应答与重发链路。

整体模块

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

后端分层

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 高并发消费链路

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 与重发

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 规划

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 约定

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 示例:

{
  "commandId": "后端下发的commandId",
  "status": "success",
  "message": "ok"
}

高并发关注点

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