- ScaleWebSocketServer:三个状态 Map 改为 ConcurrentHashMap,connectionCount 改为 AtomicInteger - ScaleServiceManager:新增 serviceScope 统一管理协程生命周期,新增 @Volatile isStarted 防重复启动 - ScaleWebSocketClient:doConnect 写入新连接前先 close() 旧连接,防止 TCP 资源泄漏 - MdnsDiscoveryManager:移除 isResolving 上多余的 @Volatile 注解 - MasterScaleActivity:全限定类名改为 import + 短类名 - WeightUtil:删除 WeightListenerImpl 中注释掉的调试代码 - scale包架构分析.md:同步更新架构文档,已修复问题压缩为汇总表 Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
248 lines
11 KiB
Markdown
248 lines
11 KiB
Markdown
# scale 包架构分析
|
||
|
||
> 最后更新:2026-05-25(同步本次优化修复内容)
|
||
|
||
## 概览
|
||
|
||
scale 包负责多设备秤数据的采集、发现、传输与聚合,采用 **mDNS + UDP 双路冗余发现 + WebSocket 长连接**的架构。
|
||
|
||
根据设备角色(MASTER/SLAVE)运行不同的服务组合:
|
||
|
||
| 类 | 子设备 | 主设备 |
|
||
|---|:---:|:---:|
|
||
| MdnsRegisterManager | ✅ | ✅ |
|
||
| ScaleWebSocketServer | ✅ | ✅ |
|
||
| UdpBroadcastSender | ✅ | ❌ |
|
||
| MdnsDiscoveryManager | ❌ | ✅ |
|
||
| UdpBroadcastReceiver | ❌ | ✅ |
|
||
| ScaleWebSocketClient | ❌ | ✅ |
|
||
| ScaleDataAggregator | ❌ | ✅ |
|
||
|
||
---
|
||
|
||
## 整体数据流
|
||
|
||
```
|
||
子设备 主设备
|
||
│ │
|
||
├─ UdpBroadcastSender ──UDP广播(8766)──▶ UdpBroadcastReceiver ─┐
|
||
├─ MdnsRegisterManager ──mDNS注册──▶ MdnsDiscoveryManager ─┤
|
||
│ │ │
|
||
│ onDeviceFound │
|
||
│ │ │
|
||
├─ ScaleWebSocketServer ◀──WS连接(8765)── ScaleWebSocketClient ◀┘
|
||
│ │ 秤数据推送 │
|
||
│ └──────────────────────────────▶ ScaleDataAggregator
|
||
│ │
|
||
│ StateFlow → UI
|
||
│
|
||
│ ◀── ScaleCommand(清零指令)─────────────────┤
|
||
│ ◀── ScaleEvent(配置同步)──────────────────┤
|
||
│ ──▶ ScaleEvent(调料添加通知)──────────────▶│
|
||
```
|
||
|
||
---
|
||
|
||
## 数据模型(3个)
|
||
|
||
### ScaleData
|
||
单个秤的数据快照,是整个包内流转的核心数据结构。
|
||
|
||
| 字段 | 类型 | 说明 |
|
||
|---|---|---|
|
||
| deviceId | String | 所属设备 ID |
|
||
| address | Int | 秤硬件地址编号 |
|
||
| weight | Double | 重量(克) |
|
||
| state | Int | 1=稳定,0=不稳定,2=量程溢出 |
|
||
| ts | Long | 数据时间戳(毫秒) |
|
||
| ip | String | 所属设备 IP,用于 UI 展示;网络传输数据中可能为空 |
|
||
| name | String? | 秤槽位名称(可选),默认为 null |
|
||
|
||
### ScaleEvent
|
||
主子设备之间的**非重量类**通知,通过 WebSocket 传输。
|
||
|
||
| 事件类型 | 方向 | 说明 |
|
||
|---|---|---|
|
||
| `seasoning_added` | 子设备 → 主设备 | 某秤检测到调料添加,携带 `delta`(重量变化量) |
|
||
| `seasoning_config` | 主设备 → 子设备 | 调料槽位配置同步,携带 `List<SlotConfig>` |
|
||
| `clear_data` | 主设备 → 子设备 | 通知所有子设备清除本机全部测试数据并重置 SP 状态 |
|
||
|
||
### ScaleCommand
|
||
主设备向子设备发送的**控制指令**,子设备收到后校验 `deviceId` 是否匹配自身再执行。
|
||
|
||
| 指令 | 说明 |
|
||
|---|---|
|
||
| `tare` | 清零指定地址的秤 |
|
||
|
||
---
|
||
|
||
## 设备发现(4个)
|
||
|
||
采用 mDNS + UDP 双路冗余,任意一路发现子设备均可触发连接。
|
||
|
||
### MdnsRegisterManager(主设备和子设备均运行)
|
||
将本机 WebSocket 服务以 mDNS 形式注册到局域网。
|
||
|
||
- 服务名格式:`DishMatch-{deviceId}`
|
||
- 服务类型:`_dishmatch._tcp.`
|
||
- 端口:`8765`
|
||
|
||
### MdnsDiscoveryManager(仅主设备运行)
|
||
持续扫描局域网中所有 `DishMatch-*` 的 mDNS 服务,解析出 IP:PORT 后触发 `onDeviceFound`。
|
||
|
||
> **关键细节**:Android `NsdManager.resolveService` 不支持并发调用,多台子设备同时被发现时会报 `FAILURE_ALREADY_ACTIVE(3)`。内部使用串行队列(`resolveQueue`)逐一解析,避免解析失败。
|
||
|
||
### UdpBroadcastSender(仅子设备运行)
|
||
每 5 秒向 `255.255.255.255:8766` 广播一个 JSON 包,作为 mDNS 的兜底发现机制。
|
||
|
||
广播包结构:
|
||
```json
|
||
{ "deviceId": "xxx", "ip": "192.168.1.x", "port": 8765 }
|
||
```
|
||
|
||
### UdpBroadcastReceiver(仅主设备运行)
|
||
监听 `8766` 端口,接收子设备的 UDP 广播包,解析后触发 `onDeviceFound`。
|
||
|
||
- 内部用 `knownDevices`(`ConcurrentHashMap`)缓存已发现的设备,避免每 5 秒重复触发连接
|
||
- 设备断线时需调用 `removeDevice()` 清除缓存,才能在重连时重新触发 `onDeviceFound`
|
||
|
||
---
|
||
|
||
## 数据传输(2个)
|
||
|
||
### ScaleWebSocketServer(主设备和子设备均运行)
|
||
基于 `java-websocket` 的服务端,监听 `8765` 端口。
|
||
|
||
**职责:**
|
||
- 监听本机 `WeightUtil` 回调,将秤数据实时推送给所有已连接客户端(节流 100ms)
|
||
- 新客户端连接时,立即推送所有秤的最新快照(`latestData` 缓存);若 `latestData` 为空(刚重启尚无读数),延迟 3 秒后补推一次
|
||
- 接收主设备下发的 `ScaleCommand`(清零)和 `ScaleEvent`(配置同步、清除数据)
|
||
- `latestData`/`lastPushTime`/`pendingPushTasks` 均使用 `ConcurrentHashMap`,保证 WeightUtil 回调线程与 java-websocket 服务端线程并发安全
|
||
- 通过 `connectionCount`(`AtomicInteger`)原子跟踪连接数,避免 `connections` 集合竞态问题
|
||
- 内置单线程 `scheduler`,用于延迟补推任务调度;`stop()` 时同步关闭
|
||
|
||
**节流策略(leading + trailing):**
|
||
1. 窗口已过(距上次推送 ≥ 100ms):立即推送,取消已有的 trailing 任务
|
||
2. 窗口内首次触发:安排一个 trailing 延迟任务,保证窗口末尾推最新值
|
||
3. 窗口内再次触发:仅更新 `latestData`,等 trailing 任务触发时推送
|
||
4. 数据未变化(weight 和 state 均相同):直接跳过,不进入节流逻辑
|
||
|
||
### ScaleWebSocketClient(仅主设备运行)
|
||
管理主设备与多台子设备的 WebSocket 长连接。
|
||
|
||
**职责:**
|
||
- 多设备并发连接(`ConcurrentHashMap` 管理)
|
||
- 断线自动重连(指数退避:2s → 4s → 8s → ... → 30s)
|
||
- 向指定设备或全部设备发送指令/事件
|
||
- 子设备首次连接成功时触发 `onDeviceConnected`,供主设备推送全量配置
|
||
|
||
**并发安全机制:**
|
||
- **版本号防重复连接**:每次调用 `connect()` 时递增 `connectVersions[deviceId]`,`doConnect` 执行前校验版本号,版本不匹配(说明已有更新的连接请求)则直接放弃,避免 UDP 触发的新连接与指数退避重连任务并发建立两条连接
|
||
- **原子移除防误删**:`onFailure`/`onClosed` 使用 `connections.remove(deviceId, webSocket)` 原子操作,只有移除的是自己的实例时才触发 `onDeviceDisconnected` 和重连,避免旧连接超时回调误删新连接引用
|
||
- **旧连接显式关闭**:`doConnect` 末尾使用 `connections.put(deviceId, ws)?.close(1000, "新连接替换")`,IP 变化重连时先关闭旧连接,防止旧 ws 对象悬空造成 TCP 资源泄漏
|
||
|
||
---
|
||
|
||
## 数据聚合(1个)
|
||
|
||
### ScaleDataAggregator(仅主设备运行)
|
||
将本机秤和所有子设备秤的数据统一汇总,以 `StateFlow` 暴露给 UI 层。
|
||
|
||
- **本机秤**:直接监听 `WeightUtil` 回调,无需经过网络
|
||
- **子设备秤**:由 `ScaleWebSocketClient.onScaleData` 回调写入
|
||
- **Map key 格式**:`{deviceId}#{address}`,便于 UI 按设备分组展示
|
||
- **IP 回填**:若数据包先于 IP 信息到达,`setDeviceIp()` 会回填已缓存数据中的空 IP 字段
|
||
|
||
---
|
||
|
||
## 配置(1个)
|
||
|
||
### ScaleDeviceConfig
|
||
硬编码各子设备的固定 UUID 和秤地址显示顺序。
|
||
|
||
| 常量 | 说明 |
|
||
|---|---|
|
||
| `DEVICE_ID_2/22/18/1` | 各子设备固定 UUID |
|
||
| `SCALE_ORDER_22` | 22个秤的物理位置排列顺序(地址范围 1-22) |
|
||
| `SCALE_ORDER_18` | 18个秤的物理位置排列顺序(地址范围 1-18) |
|
||
| `DEVICE_ORDER` | 设备在列表中的显示顺序 |
|
||
|
||
---
|
||
|
||
## 门面(1个)
|
||
|
||
### ScaleServiceManager
|
||
整个包的统一入口(单例),根据设备角色决定启动哪些服务,并将各组件串联起来。
|
||
|
||
- **防重复启动**:内置 `@Volatile isStarted` 标志,`start()` 入口加守卫,重复调用直接返回,防止端口冲突和资源泄漏
|
||
- **协程生命周期管理**:维护 `serviceScope`(`Dispatchers.IO + SupervisorJob()`),所有内部异步操作均在此 Scope 内执行,`stop()` 时统一 `cancel()`,防止 NPE 和内存泄漏
|
||
|
||
**外部使用方式:**
|
||
```kotlin
|
||
// Application.onCreate()
|
||
ScaleServiceManager.start(context)
|
||
|
||
// 主设备 UI 订阅全量秤数据
|
||
ScaleServiceManager.allScales?.collect { scales -> ... }
|
||
|
||
// 主设备发送清零指令
|
||
ScaleServiceManager.sendTare(deviceId, address)
|
||
|
||
// 主设备广播调料配置
|
||
ScaleServiceManager.sendSeasoningConfig(slots)
|
||
|
||
// 主设备通知所有子设备清除数据
|
||
ScaleServiceManager.sendClearData()
|
||
|
||
// 主设备监听子设备秤事件(如调料添加)
|
||
ScaleServiceManager.onScaleEvent = { event -> ... }
|
||
|
||
// 子设备监听主设备下发的调料配置同步
|
||
ScaleServiceManager.onSeasoningConfig = { event -> ... }
|
||
|
||
// 子设备监听主设备下发的清除数据指令
|
||
ScaleServiceManager.onClearData = { ... }
|
||
|
||
// 子设备向主设备广播秤事件
|
||
ScaleServiceManager.broadcastEvent(event)
|
||
|
||
// 子设备监听主设备连接状态变化
|
||
ScaleServiceManager.onMasterConnectionChanged = { connected -> ... }
|
||
|
||
// 子设备查询当前是否有主设备连接
|
||
ScaleServiceManager.isMasterConnected
|
||
|
||
// Application.onTerminate()
|
||
ScaleServiceManager.stop()
|
||
```
|
||
|
||
外部代码只需与 `ScaleServiceManager` 交互,无需感知内部任何组件。
|
||
|
||
---
|
||
|
||
## 问题排查与优化建议
|
||
|
||
### ✅ 已修复问题(2026-05-25)
|
||
|
||
| 严重度 | 问题 | 位置 | 修复方式 |
|
||
|---|---|---|---|
|
||
| 🔴 严重 | 三个状态 Map 非线程安全 | `ScaleWebSocketServer.kt` | 改为 `ConcurrentHashMap` |
|
||
| 🔴 严重 | `connectionCount` 非原子操作 | `ScaleWebSocketServer.kt` `InternalServer` | 改为 `AtomicInteger` |
|
||
| 🟡 中等 | `CoroutineScope` 无生命周期绑定 | `ScaleServiceManager.kt` | 改为 `serviceScope`,`stop()` 时统一 cancel |
|
||
| 🟡 中等 | `start()` 缺少防重复启动保护 | `ScaleServiceManager.kt` | 新增 `@Volatile isStarted` 守卫 |
|
||
| 🟡 中等 | `doConnect` 旧连接未显式关闭 | `ScaleWebSocketClient.kt` | 改为 `connections.put()?.close()` |
|
||
| 🟡 中等 | `@Volatile` 注解多余 | `MdnsDiscoveryManager.kt` | 移除 `@Volatile` |
|
||
| 🟡 中等 | 全限定类名违反导入规范 | `MasterScaleActivity.kt` | 改为 import + 短类名 |
|
||
| 🟢 轻微 | `onGetWeight` 残留注释掉的调试代码 | `WeightUtil.kt` | 删除注释代码块 |
|
||
| 🟢 轻微 | "方案一"/"方案三"草稿注释 | `ScaleWebSocketServer.kt` | 改为描述行为意图的正式注释 |
|
||
|
||
---
|
||
|
||
### 🟢 轻微:ScaleDataAggregator 高频更新时 GC 压力(待观察)
|
||
|
||
**位置**:`ScaleDataAggregator.kt`,第 120 行
|
||
|
||
**问题**:每次秤数据更新都执行 `_allScales.value = HashMap(cache)`,在多台设备同时高频推送数据时(如 22 个秤 × 10Hz),每秒可能创建数百个临时 HashMap 对象,增加 GC 压力。
|
||
|
||
**建议**:这是 StateFlow 的惯用写法,短期内无需优化。若未来出现 GC 卡顿,可考虑引入防抖(debounce)合并多次更新后再发布快照。
|