搜索页最让人误判的故障,不是请求报错,而是页面看起来还能用。用户连续输入四次,最后一个关键词已经显示在输入框里,列表却在半秒后跳回第二次查询;从详情页返回后,同一次输入又触发两遍加载。日志没有崩溃,接口也都是 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 布尔值猜测生命周期。

type SearchState = 'IDLE' | 'TYPING' | 'LOADING' |
  'READY' | 'ERROR' | 'DISPOSED'

interface SearchItem {
  id: string
  title: string
}

interface SearchSnapshot {
  runId: string
  requestSeq: number
  generation: number
  keyword: string
  state: SearchState
  requests: number
  committed: number
  stale: number
  results: SearchItem[]
}

const initialSnapshot: 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'

interface SearchEnvelope {
  seq: number
  generation: number
  keyword: string
  items: SearchItem[]
}

class SearchStreamController {
  private query$ = new Subject<string>()
  private dispose$ = new Subject<void>()
  private subscription?: Subscription
  private seq: number = 72
  private generation: number = 18
  snapshot: SearchSnapshot = { ...initialSnapshot }

  constructor(private searchApi: (q: string) => Promise<SearchItem[]>) {}

  start(): void {
    if (this.subscription && !this.subscription.closed) return

    this.subscription = this.query$.pipe(
      map(value => value.trim()),
      debounceTime(180),
      distinctUntilChanged(),
      switchMap(keyword => {
        const seq = ++this.seq
        const generation = this.generation
        this.snapshot = {
          ...this.snapshot, keyword, requestSeq: seq,
          requests: this.snapshot.requests + 1, state: 'LOADING'
        }
        return defer(() => 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') return
    this.snapshot = { ...this.snapshot, state: 'TYPING' }
    this.query$.next(value)
  }

  private commit(result: SearchEnvelope): void {
    const current = result.generation === this.generation &&
      result.seq === this.seq
    if (!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'
    }
  }

  private fail(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、注销订阅并令旧页面失去提交资格,防止返回页面后订阅数累加。

class SearchStreamController {
  // 省略前文成员与 start/input

  dispose(): void {
    if (this.snapshot.state === 'DISPOSED') return
    this.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 = undefined
    console.info('[RX-2911] activeSubscriptions=0 state=DISPOSED')
  }
}

@Entry
@Component
struct ReactiveSearchPage {
  private controller = new SearchStreamController(queryCatalog)
  @State keyword: 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 和状态仓储必须隔离;“全局永远只能有一条”并不是本文结论。本文真正要求的是每个所有者有明确预算,并在自己的终点成对释放。

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

官方与上游参考:

Logo

作为“人工智能6S店”的官方数字引擎,为AI开发者与企业提供一个覆盖软硬件全栈、一站式门户。

更多推荐