- 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>
11 KiB
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 的兜底发现机制。
广播包结构:
{ "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):
- 窗口已过(距上次推送 ≥ 100ms):立即推送,取消已有的 trailing 任务
- 窗口内首次触发:安排一个 trailing 延迟任务,保证窗口末尾推最新值
- 窗口内再次触发:仅更新
latestData,等 trailing 任务触发时推送 - 数据未变化(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 和内存泄漏
外部使用方式:
// 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)合并多次更新后再发布快照。