
一、问题从一条重复的派单开始OpsStream 是配送运维看板页面LiveOpsPage订阅dispatch-east频道把调度事件实时更新到列表。测试同事说它“切后台再回来会偶尔多一条相同任务”。服务端查了发送记录每个 seq 只有一份客户端 HiLog 却在 11:42 连续打印两次seq18417。当时页面顶部仍写着“连接正常”所以这个问题一度被当成列表 diff 错误。沿着订阅链倒查才发现页面aboutToAppear()每次都会调用connect()WebSocket 的旧message监听还在RxJS 上层又新建了一次订阅。更糟的是断网期间服务端产生了 6 条事件重连后客户端只从“当前最新”继续重复和丢失竟然同时存在。任务WS-1142最后不是靠列表去重补丁解决而是重新定义了连接、回放和 UI 订阅各自的边界。二、把“连上了”拆成四个可观测状态HarmonyOS 的 WebSocket 流程很明确createWebSocket()创建对象on(open|message|close|error)订阅事件connect()建连结束时close()并配合off()移除监听。问题在于系统 API 只告诉我们传输层发生了什么产品还需要知道回放是否完成。因此 OpsStream 不再使用一个connected: boolean而是四态DISCONNECTED → CONNECTING → SYNCING → LIVE。open只代表物理连接建立此时发送resumeFromlastAck1服务端回放结束帧到达后才进入LIVE。如果第一条实时事件的 seq 大于期望值就退回SYNCING请求缺口而不是直接渲染一个看似正常、实际断档的列表。最终运行值固定下来会话ws_72最后确认序号 18420重连从 18421 请求恢复缺口 6 条丢弃重复 4 条活跃订阅者 1 个重连 2 次端到端 P95 为 88 ms。这些值同时出现在状态仓、HiLog 和手机页不再各算各的。三、目录按“传输、序列、视图”切开工程名OpsStream页面没有持有原生 WebSocket 对象。项目结构如下entry/src/main/ets/ ├── pages/LiveOpsPage.ets ├── stream/HarmonySocketSource.ets ├── stream/DispatchStreamStore.ets ├── stream/SequenceGate.ets └── model/DispatchEvent.etsHarmonySocketSource只负责把回调 API 转成 ObservableSequenceGate维护expectedSeq、重复计数和缺口DispatchStreamStore负责重连、回放状态与共享页面只拿StreamSnapshot。把三层分开之后测试可以向SequenceGate注入乱序数组不必真的断一次网。四、第一段代码原生回调只对应一个 Observable 生命周期之前的适配器在connect()里不断on(message)却只在应用退出时关闭。新的源 Observable 在订阅时创建 WebSocket在 teardown 中成对执行off()和close()。下面省略证书与鉴权头只保留 Demo 真正参与的生命周期。import{webSocket}fromkit.NetworkKit;import{BusinessError}fromkit.BasicServicesKit;import{Observable}fromrxjs;exportfunctionsocketFrames(url:string):Observablestring{returnnewObservablestring((subscriber){constwswebSocket.createWebSocket();constonOpen():voidsubscriber.next({type:OPEN});constonMessage(_err:BusinessError,value:string|ArrayBuffer):void{if(typeofvaluestring)subscriber.next(value);};constonClose():voidsubscriber.complete();constonError(err:BusinessError):voidsubscriber.error(err);ws.on(open,onOpen);ws.on(message,onMessage);ws.on(close,onClose);ws.on(error,onError);ws.connect(url);return(){ws.off(open,onOpen);ws.off(message,onMessage);ws.off(close,onClose);ws.off(error,onError);ws.close();};});}这里故意保留四个稳定函数引用因为off()必须拿到与on()对应的回调。teardown 不是只在页面销毁时执行switchMap切换连接、重试替换源、最后一个订阅者离开时都会触发。close()可能在未完成建连时调用所以工程版会把同步异常记录为closeSkipped但不会让 teardown 再向已结束的 subscriber 抛错。另一个边界是二进制帧。Demo 协议只接受 JSON 文本遇到ArrayBuffer直接计入协议错误而不是用默认编码猜测。若产品要传 protobuf应在这一层明确解码并把解码失败与网络断开分开统计。五、第二段代码序号闸门决定重复、缺口和确认点列表组件做id去重太晚了因为重复事件可能已经触发通知、振动或统计。SequenceGate在任何副作用之前检查 seq。它只接受三种结果ACCEPT、DUPLICATE、GAP只有连续事件才能推进lastAck。exporttypeGateResultACCEPT|DUPLICATE|GAP;exportclassSequenceGate{lastAck:number18420;duplicates:number0;gapSize:number0;inspect(seq:number):GateResult{if(seqthis.lastAck){this.duplicates;returnDUPLICATE;}constexpectedthis.lastAck1;if(seqexpected){this.gapSizeseq-expected;returnGAP;}this.lastAckseq;this.gapSize0;returnACCEPT;}}为什么不接受“后到的大序号再等小序号补齐”因为 OpsStream 的协议保证单频道单调发送出现 gap 就说明客户端确实缺数据。保持一个小型乱序缓存只会掩盖服务端或链路问题。lastAck在事件完成业务落库后更新而不是一收到帧就更新否则渲染过程崩溃重连也不会再要这条数据。本轮注入测试先收到 18423再回放 18421、18422、18423。第一帧触发缺口 2回放推进确认点第二次出现的 18423 被判为重复。最终 6 条缺口全部恢复重复计数累计 4没有对列表产生二次副作用。六、第三段代码共享连接但不把缓存变成僵尸RxJS 的shareReplay默认refCount为 false这意味着订阅者归零后源可能仍保持订阅。对 HTTP 缓存或许有用对页面级 WebSocket 却会留下后台连接。Store 使用shareReplay({ bufferSize: 1, refCount: true })让状态栏和列表共享同一条源同时保证最后一个消费者离开后触发原生 teardown。import{BehaviorSubject,Observable,timer}fromrxjs;import{retry,shareReplay,switchMap}fromrxjs/operators;exportclassDispatchStreamStore{privatereadonlychannel$newBehaviorSubjectstring(dispatch-east);readonlyevents$:ObservableDispatchEventthis.channel$.pipe(switchMap((channel)socketFrames(this.url(channel)).pipe(retry({count:3,delay:(_error,retryCount)timer(retryCount*800)}),this.decodeAndResume(channel))),shareReplay({bufferSize:1,refCount:true}));setChannel(channel:string):void{if(channel!this.channel$.value)this.channel$.next(channel);}dispose():void{this.channel$.complete();}}switchMap解决频道切换新频道进入时旧 Observable 先 teardown旧 socket 的晚帧没有机会混进新页面。retry放在共享之前所有消费者共同经历一次重连如果放在页面各自的订阅之后状态栏和列表会各自建立一条连接。800、1600、2400 ms 的有限退避只处理可恢复网络错误鉴权失败不会重试。bufferSize: 1缓存的是最近一条已校验事件不是完整历史。页面重新订阅能立即得到当前态但仍要通过resumeFrom补齐离线区间。离开页面时两个 UI 订阅都取消引用计数归零原生回调与 socket 随之释放aboutToDisappear()另外调用store.dispose()防止页面实例留在导航缓存时继续接受频道切换。调试阶段我增加了subscribers计数并让每次源订阅打印sessionId。修复前回前台一次会看到ws_70、ws_71同时接收修复后只有ws_72。关键日志为taskWS-1142 sessionws_72 stateLIVE lastAck18420 subscribers1下一行是resumeFrom18421 gapRecovered6 duplicates4 reconnects2 p9588ms。七、恢复边界比重连次数更重要把 socket 自动连回来只是第一层。真正决定数据是否可信的是回放窗口。服务端为dispatch-east保留最近 10 分钟或 5000 条事件客户端若携带的resumeFrom早于窗口下界会收到RESET_REQUIRED。此时不能假装补齐成功而要清空增量状态调用一次快照接口再以快照的headSeq建立实时流。相反短暂断网不应该全量刷新。lastAck18420时只请求 18421 之后的事件收到 6 条回放并见到REPLAY_END后进入LIVE。回放期间 UI 仍显示旧列表但顶部明确标为“同步中”避免用户把未完成的列表当成实时结果。前后台也是一个边界。进入后台后 OpsStream 主动退订而不是依赖系统最终清理回前台重新订阅共享源创建一个新 socket并从持久化的确认点恢复。确认点每 20 条或 2 秒批量写入 Preferences减少频繁 I/O。若进程在批量写入前退出最多重复回放 20 条序号闸门能安全丢弃而不会丢数据。八、最终运行状态能互相对账在飞行模式开关、前后台切换、频道切换三组测试中当前会话稳定为ws_72状态从CONNECTING经过SYNCING回到LIVE。最后确认序号 18420离线缺口恢复 6 条由历史确认点和重试产生的 4 条重复全部在副作用前丢弃活跃订阅者始终为 1。两次重连的消息端到端 P95 是 88 ms。手机页 11:42 展示同一组数据并把resumeFrom 18421放在“回放边界”卡片里。按钮只有“切换频道”和“重新同步”没有虚构一个“修复网络”的营销动作。用户看到的是当前数据可信度而不只是绿色连接图标。九、这次留下的不是一个重试操作符最终方案的核心并非 RxJS 写法而是三个所有权原生监听属于源 Observable 的订阅周期序号确认点属于业务副作用完成之后共享连接属于仍在场的消费者。三个边界一旦混在页面生命周期里双订阅、僵尸连接和错误确认迟早会同时出现。这个结构也有适用范围。若频道事件不保证单调SequenceGate必须改成带窗口的重排器若每条事件都要求强事务确认Preferences 批量持久化不够需要关系型存储若后台必须持续接收则不能简单依赖页面refCount应把 Store 上移到明确的长生命周期能力中。工程封装不是把所有情况藏起来而是把当前协议的假设写得足够清楚。参考资料华为开发者文档《WebSocket Connection》页面更新于 2025-05-20https://developer.huawei.com/consumer/en/doc/harmonyos-guides-V14/websocket-connection-V14 RxJSshareReplayAPIhttps://rxjs.dev/api/index/function/shareReplay