☰
HarmonyOS 7 + RxJS + ArkUI:搜索事件流的晚到结果隔离与订阅收口【鸿蒙心迹】
2026/10/3 12:40:58 网站建设 项目流程

搜索页最让人误判的故障,不是请求报错,而是页面看起来还能用。用户连续输入四次,最后一个关键词已经显示在输入框里,列表却在半秒后跳回第二次查询;从详情页返回后,同一次输入又触发两遍加载。日志没有崩溃,接口也都是 200,问题藏在事件流的所有权里。

本文用 StreamShelf 演示一个边界明确的改造:ArkUI 负责输入和状态呈现,RxJS 负责事件组合,业务层用 runId 与 generation 限制最终提交。演示任务号为 RX-2911,页面为 ReactiveSearchPage,时间为 19:18。四次输入产生四个请求,只有最后一个结果可以提交;页面退出后活动订阅由 1 归零。

一、先看结果,再拆开那三次“无效成功”

演示查询依次是“折”“折叠”“折叠屏”“折叠屏适配”。网络层记录 requestSeq 73、74、75、76,响应顺序却是 74、73、76、75。若页面拿到结果就直接赋值,requestSeq 75 最晚返回,会覆盖正确的 76。

StreamShelf 的最终摘要是:requests=4、committed=1、stale=3、results=12、duration=268 ms。状态从 IDLE 进入 TYPING,再到 LOADING、READY;离开页面后进入 DISPOSED。这里的 268 ms、12 条结果和三次丢弃都是演示数据,用于把正文、日志和图片对齐,不代表某个真实服务的性能。

引入 RxJS 的目的也不是把 Promise 换一种写法。真正需要解决的是三个问题:输入事件如何合并,旧查询的结果如何失去提交资格,页面销毁时整条订阅链如何结束。如果只写 debounceTime,没有处理后两个问题,重复订阅和晚到结果仍会出现。

二、switchMap 能切换订阅,不能凭空取消所有工作

RxJS 官方把 Observable 定义为组织异步与事件程序的核心类型。switchMap 在新的外层值到来时切换到新的内部 Observable,这很适合搜索输入。但一个常被省略的事实是:如果内部工作只是由普通 Promise 包装,取消订阅不等于底层网络请求一定中止。

因此文章不把 switchMap 描述成“自动取消接口”。它能阻止旧内部流继续向下游发值;真正的传输取消要看请求库是否提供 AbortSignal、cancel token 或对应能力。没有传输取消时,旧 Promise 仍可能完成,只是不能再改页面状态。业务层仍需 generation 作为最后一道提交门禁。

另一个误区是把 Subject 做成应用全局单例。页面每次进入都调用 subscribe,却没有对应 unsubscribe,订阅数会随访问次数增加。输入一次,多个旧页面实例同时处理;即使这些实例已经不可见,它们仍可能写日志、请求网络或持有闭包引用。

三、先把状态写成可核对的合同

这段代码解决什么问题:把查询代次、页面状态和演示指标集中到一个模型里,避免 UI 通过 loading 布尔值猜测生命周期。

typeSearchState='IDLE'|'TYPING'|'LOADING'|'READY'|'ERROR'|'DISPOSED'interfaceSearchItem{id:stringtitle:string}interfaceSearchSnapshot{runId:stringrequestSeq:numbergeneration:numberkeyword:stringstate:SearchState requests:numbercommitted:numberstale:numberresults:SearchItem[]}constinitialSnapshot:SearchSnapshot={runId:'RX-2911',requestSeq:72,generation:18,keyword:'',state:'IDLE',requests:0,committed:0,stale:0,results:[]}

requestSeq 用于区分一次运行里的请求,generation 用于区分页面实例。它们不是同一个概念。用户每次输入会增加 requestSeq;页面重新进入时 generation 改变,即使 requestSeq 恰好相同,旧页面也没有提交权。

状态里保留 requests、committed 和 stale,是为了让诊断页能解释“为什么发了四次却只更新一次”。实际产品不一定把这些字段展示给用户,但日志和测试断言应该能读取。results 只在通过提交门禁后替换,不能在网络回调里先清空再判断,否则旧请求仍会制造一次视觉闪烁。

四、把输入流、请求流和提交点放在同一条管线里

这段代码解决什么问题:对输入做去抖与去重,用 switchMap 切换内部流,并在 Promise 无法中止时用 generation 和 requestSeq 拒绝晚到结果。

import{Subject,Subscription,debounceTime,distinctUntilChanged,switchMap,takeUntil,defer,from,map,finalize}from'rxjs'interfaceSearchEnvelope{seq:numbergeneration:numberkeyword:stringitems:SearchItem[]}classSearchStreamController{privatequery$=newSubject<string>()privatedispose$=newSubject<void>()privatesubscription?:Subscriptionprivateseq:number=72privategeneration:number=18snapshot:SearchSnapshot={...initialSnapshot}constructor(privatesearchApi:(q:string)=>Promise<SearchItem[]>){}start():void{if(this.subscription&&!this.subscription.closed)returnthis.subscription=this.query$.pipe(map(value=>value.trim()),debounceTime(180),distinctUntilChanged(),switchMap(keyword=>{constseq=++this.seqconstgeneration=this.generationthis.snapshot={...this.snapshot,keyword,requestSeq:seq,requests:this.snapshot.requests+1,state:'LOADING'}returndefer(()=>from(this.searchApi(keyword))).pipe(map(items=>({seq,generation,keyword,items})),finalize(()=>console.info(`[RX-2911] seq=${seq}finalized`)))}),takeUntil(this.dispose$)).subscribe({next:(result:SearchEnvelope)=>this.commit(result),error:(error:Error)=>this.fail(error)})}input(value:string):void{if(this.snapshot.state==='DISPOSED')returnthis.snapshot={...this.snapshot,state:'TYPING'}this.query$.next(value)}privatecommit(result:SearchEnvelope):void{constcurrent=result.generation===this.generation&&result.seq===this.seqif(!current){this.snapshot={...this.snapshot,stale:this.snapshot.stale+1}return}this.snapshot={...this.snapshot,results:result.items,committed:this.snapshot.committed+1,state:'READY'}}privatefail(error:Error):void{console.error(`[RX-2911]${error.message}`)this.snapshot={...this.snapshot,state:'ERROR'}}}

180 ms 是 Demo 参数,不是通用最佳值。它要根据输入法、接口成本和交互目标测量。distinctUntilChanged 只会过滤连续相同字符串;若业务把大小写、全角字符或别名视为等价,还需要在它之前做规范化。

finalize 会在内部流完成、报错或被取消订阅时执行,适合回收这次内部流的计数、追踪或局部资源,但不应该在里面直接把整个页面改成 READY。switchMap 切走旧请求时也会触发旧内部流 finalize,若 finalize 无条件关闭 loading,新的请求刚开始就可能被旧请求关掉。

图中的 DevEco Studio 画面是本批生成的演示配图,不是实机或真实 IDE 运行证据。目录、代码、任务号和 HiLog 均按本文状态模型统一。

五、错误流如果终止,搜索框会“看着能输但不再工作”

上面的简化代码把 error 放到订阅终点,便于说明错误状态,但产品实现通常不希望一次 500 让整条查询流永久终止。更合适的做法是在 switchMap 的内部流使用 catchError,把本次失败映射为可展示结果,外层 query$ 继续存活。

错误恢复也不能复用旧 seq。用户点击重试时,应产生新的请求序号,并记录 retryOf=75 之类的关联字段。否则旧请求晚到和重试结果会共享身份,日志无法判断究竟哪一次取得提交权。

网络层若支持真正取消,要把取消句柄绑定到内部 Observable 的 teardown。即便如此,提交门禁仍不能删:取消可能发生在响应已经抵达之后,缓存层也可能同步返回,页面 generation 仍是跨实例隔离的必要条件。

六、页面退出必须有一个唯一、幂等的 dispose

这段代码解决什么问题:结束 Subject、注销订阅并令旧页面失去提交资格,防止返回页面后订阅数累加。

classSearchStreamController{// 省略前文成员与 start/inputdispose():void{if(this.snapshot.state==='DISPOSED')returnthis.generation++this.snapshot={...this.snapshot,generation:this.generation,state:'DISPOSED',results:[]}this.dispose$.next()this.dispose$.complete()this.query$.complete()this.subscription?.unsubscribe()this.subscription=undefinedconsole.info('[RX-2911] activeSubscriptions=0 state=DISPOSED')}}@Entry@Componentstruct ReactiveSearchPage{privatecontroller=newSearchStreamController(queryCatalog)@Statekeyword:string=''aboutToAppear():void{this.controller.start()}aboutToDisappear():void{this.controller.dispose()}}

takeUntil 和显式 unsubscribe 看起来重复,实际承担的角色不同:dispose$ 让管线内部按声明式路径结束,subscription.unsubscribe 是对象级兜底。两者都必须幂等。更重要的是 Subject 已 complete 后不能用于新页面,所以 controller 应跟随页面实例重建,而不是 dispose 后再次 start。

aboutToDisappear 是否等同最终销毁,要结合页面和导航设计确认。有些场景页面暂时不可见但会复用,团队可以选择在不可见时暂停输入流、在最终销毁时 complete。文章采用“离开即释放”的资料搜索页策略,不应机械复制到需要保留长连接的页面。

七、手机运行图只证明状态可解释

19:18 的手机运行页展示 query=折叠屏适配、requestSeq=76、generation=18、state=READY。摘要为 4 requests / 1 committed / 3 stale,结果数 12。红色标注指向最终提交序号,说明列表来自最新查询,而不是按响应到达顺序更新。

详情诊断页则展示响应顺序 74→73→76→75,以及 seq 73、74、75 被标记 STALE_DROPPED。退出页面后 activeSubscriptions 从 1 变成 0,状态进入 DISPOSED。两张手机图承担不同任务:一张呈现用户可见结果,另一张解释为什么旧结果没有提交。

这些画面没有证明 RxJS 在所有 ArkTS 工程中都可直接采用。ArkTS 支持与 TS/JS 生态互操作,但具体三方包仍需在目标 SDK、编译模式和依赖管理环境中验证。若团队的静态检查、包体或性能要求不接受 RxJS,同样的状态机也可以用原生 Promise、定时器和序号实现。

八、调试时别只盯请求数量

第一组用例连续输入四次并刻意打乱响应,断言最终 seq=76。第二组在 LOADING 时离开页面,断言之后没有状态提交。第三组进入、退出页面五次,每次输入一次,活动订阅峰值始终为 1,结束为 0。第四组让请求报错后继续输入,验证错误是否意外终止外层流。

还要检查空关键词。trim 后为空时,可以切换到 EMPTY 并返回空列表,不应继续调用服务。中文输入法组合阶段是否每次触发 onChange,需要按真实输入行为验证;如果产品要求用户点击搜索再执行,就不该为了使用 RxJS 强行改成实时查询。

性能日志不记录完整搜索内容时,可以保存 keywordLength 与脱敏哈希。本文为了图文一致直接展示“折叠屏适配”,生产环境要根据业务数据敏感度决定。日志至少包含 runId、seq、generation、状态迁移、耗时、提交或丢弃原因。

九、给每一条订阅标出所有者

排查重复订阅时,我更愿意先画一张“谁创建、谁释放”的清单,而不是全局搜索 subscribe。页面输入订阅属于页面实例;登录态或网络状态订阅可能属于 UIAbility;应用级事件才有资格跟随进程。生命周期不同的订阅放进同一个 CompositeSubscription 或数组里,最终仍会遇到某个页面退出时误杀全局流,或者全局容器一直留着页面闭包。

StreamShelf 为每次 start 生成 ownerTag=ReactiveSearchPage#g18,HiLog 在订阅建立时写 active=1,dispose 后写 active=0。如果 start 被重复调用,方法直接返回并记录 duplicate_start,而不是再创建第二条管线。这个小字段比“应该只调用一次”的注释可靠,因为自动化测试能断言它。

订阅所有者还决定错误展示。页面流失败可以在页面显示重试;UIAbility 级网络监控失败不应该直接把某个页面切成 ERROR。把所有错误都汇入一个全局 Subject,看起来统一,实际丢掉了恢复边界。事件可以共享,状态提交仍要回到明确所有者。

十、shareReplay 不是免费的页面缓存

搜索页常见的下一步是加入 shareReplay,让多个组件复用同一结果。这个操作符很方便,也很容易把“最近一次值”变成“旧页面仍被引用”。如果上游永不完成、订阅计数策略不清楚,缓存可能比页面活得更久,连同结果数组和闭包一起保留。

因此缓存策略要先回答三个问题:缓存键是什么,缓存多久,最后一个订阅离开后是否清理。关键词相同但用户、语言、筛选条件或数据版本不同,不能共用一份缓存。页面级复用可以把缓存放在 controller 内,dispose 时整体释放;跨页面缓存则应成为单独仓储,由容量、TTL 和账号切换规则管理,不再假装只是一个操作符。

本文没有在主代码中加入 shareReplay,就是为了让生命周期可见。先证明活动订阅能从 1 回到 0,再决定是否缓存。否则一次性能优化可能同时引入数据串号和内存保留,日志只看到请求减少,却看不到旧结果的所有权变模糊。

十一、真正取消请求要给适配器 teardown 能力

如果 Network Kit 或业务请求封装能够暴露取消句柄,可以用 Observable 构造函数把它绑定到 teardown:订阅时创建请求,响应时 next/complete,取消订阅时调用 cancel。这样 switchMap 切换旧内部流时,传输层也收到停止信号。这里的重点不是照搬某个网络库的名字,而是让适配器显式声明“可取消”还是“只能忽略结果”。

两种适配器必须用不同指标。可取消请求统计 canceledTransport;不可取消 Promise 统计 staleResultDropped。把二者都写成 canceled 会夸大优化效果。服务端已经收到并处理的请求仍会消耗资源,客户端只是没有提交结果。

取消还要处理竞态:响应和 teardown 可能几乎同时发生。适配器需要一个 settled 标记,确保 complete、error 和 cancel 只有一个成为终态。业务 generation 仍保留,因为页面退出与账号切换属于更高层的权限变化,不该依赖底层请求是否来得及取消。

十二、ArkUI 状态更新要以快照提交

ReactiveSearchPage 不应分别更新 loading、items、errorText、requestSeq 四个 @State 字段。分散赋值会在某个渲染帧出现“新关键词配旧列表”或“READY 但 loading 仍为 true”的短暂组合。controller 应生成不可变快照,页面一次替换视图模型,让字段来自同一个 seq。

快照也让详情图更可信。诊断页看到 seq=76、state=READY、results=12 时,三者来自一次 commit;如果逐字段更新,截图可能刚好落在中间态。实际项目可通过 @ObservedV2、状态管理方案或仓储适配完成,关键不是选哪个装饰器,而是保持原子提交语义。

列表渲染还要使用稳定业务 ID。晚到结果被拒绝解决了顺序问题,但如果每次结果都用数组下标作为 key,列表仍可能复用错误节点,表现成标题和缩略图短暂错位。异步正确性与渲染身份是两条边界,都要分别测试。

十三、引入第三方库之前先做四项验收

第一项是编译兼容:在目标 HarmonyOS SDK、ArkTS 模式和 DevEco Studio 版本下构建最小 Demo,不根据“TypeScript 库”三个字推断必然可用。第二项是包体与依赖:记录引入前后的产物差异,检查是否带入不需要的模块。第三项是运行语义:验证定时器、Promise、错误栈和取消路径符合预期。第四项是维护边界:锁定版本、记录上游来源和许可证,避免下次安装得到不同实现。

如果只用 debounce、switchLatest 和 dispose 三个概念,自研几十行协调器可能更透明;如果页面需要合并输入、网络、缓存、筛选、重试和生命周期,RxJS 的操作符组合会更有价值。评估不该围绕“响应式更高级”,而应比较团队能否读懂错误路径、是否能写出确定测试、依赖成本是否接受。

StreamShelf 还准备了一个无 RxJS 的对照实现,保持同样的 runId、seq 和 generation。两套实现跑相同乱序用例,结果都必须是 4/1/3。只有语义一致后,才讨论哪套代码更容易维护。这样三方库是实现选择,不会成为业务合同本身。

十四、测试时间要可控,不能真的等网络运气

乱序测试如果依赖真实接口延迟,每次结果都可能不同,无法稳定覆盖 seq 73、74、75 被拒绝的分支。测试替身应接收关键词与序号,由用例明确安排 resolve 顺序。先发四个请求,再按 74、73、76、75 完成,最后断言列表来自 76、committed=1、stale=3。

debounceTime 也不应让测试真的等待 180 ms。可以把时钟或调度策略作为 controller 的构造依赖,在生产使用真实时间,在单元测试推进虚拟时间。文章示例为了阅读简化了这个入口,项目落地时应避免把毫秒常量散落在操作符链和测试代码里。

页面退出用例要在 Promise 完成前调用 dispose,再完成旧 Promise,断言快照仍为 DISPOSED。这个顺序专门验证“取消订阅不等于 Promise 消失”的边界。若只在请求完成后退出,最关键的晚到回调分支从未被执行,覆盖率数字再高也没有意义。

十五、从日志反推问题时按三层排查

第一层看输入:是否产生预期 keyword,规范化后是否为空,debounce 是否合并。第二层看流:seq 是否单调、switchMap 是否 finalize 旧内部流、错误是否终止外层。第三层看提交:generation、seq 和页面状态是否同时满足。如果输入只有一次却请求两次,优先查重复订阅;请求四次但提交旧结果,查门禁;提交正确却列表错位,查 ArkUI key。

这套顺序能避免把所有问题都归因于网络。HiLog 中[RX-2911] seq=75 finalized只说明内部流结束,不说明传输取消;STALE_DROPPED说明结果到达但无提交权;activeSubscriptions=0说明页面资源已收口。三个日志分别对应事件、结果与生命周期,不能互相替代。

在评审里还应要求失败样本附上状态时间线,而不是只贴最终截图。截图能展示 4/1/3,却不能证明 75 为什么晚到。时间线与快照结合,才足以定位责任层。

十六、适用边界与最终判断

RxJS 适合多个异步事件需要组合、切换和统一释放的页面。只有一次按钮请求的简单页面,引入整套流库可能增加理解与包体成本。switchMap 解决的是内部订阅切换,不是对任意 Promise 的魔法取消;takeUntil 解决的是管线结束,不替代页面实例的所有权设计。

StreamShelf 最终把问题拆成三层:ArkUI 产生输入意图,RxJS 组合事件,generation 与 requestSeq 决定提交资格,dispose 收回页面资源。四次请求只提交一次并不代表浪费已经完全消失,却保证旧结果不会污染当前页面。工程上先把正确性和生命周期做清楚,再评估是否需要真正的传输取消与缓存,顺序更稳。

还有一个容易漏掉的验收点:热重载、预览器和测试框架可能让页面生命周期与正式设备不同。开发阶段看到 activeSubscriptions=0 仍要在目标运行形态复测,尤其是 Navigation 缓存页面、Tabs 保活和多窗口同时打开的情况。若产品允许两个页面实例并存,它们可以各有一条订阅,但 ownerTag、generation 和状态仓储必须隔离;“全局永远只能有一条”并不是本文结论。本文真正要求的是每个所有者有明确预算,并在自己的终点成对释放。

只有边界和验收同时成立,操作符组合才真正服务于工程,而不是增加一层难以追踪的语法。

官方与上游参考:

  • ArkTS 与 TypeScript/JavaScript 生态兼容说明
  • RxJS 概览
  • RxJS switchMap
  • RxJS takeUntil
  • RxJS API 索引与 finalize

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询