You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJS中mergeMap置于switchMap内与后置的执行差异问题

RxJS两段groupBy代码行为差异的根本原因

两者行为差异的核心是流的生命周期边界不同,本质是对toArray、groupBy两个操作符的触发规则、以及switchMap的完成时机理解有偏差,先明确两个关键操作符的特性:

  • toArray():只有当它订阅的上游Observable发送complete(完成)信号时,才会把期间收集到的所有元素打包成数组发射;只要上游不完成,它会一直缓存元素,永远不会输出值。
  • groupBy(keySelector):只要它的上游不发送完成信号,它创建的每个分组Observable就会一直保持存活,持续等待后续匹配对应key的元素进入,不会主动触发完成。

第一段代码无输出的原因

第一段代码的管道结构如下:

const arraySubject = new Subject();

arraySubject.pipe(
  switchMap((arr) => from(arr)),
  groupBy((elem) => elem.type),
  mergeMap((group) => group.pipe(toArray()))
).subscribe((arrGroups) => {
  // nothing gets logged
})

arraySubject.next([{name: 'Bill', type: 'user'}, {name: 'Jane', type: 'admin'}, {name: 'John', type: 'user'}, {name: 'Ann', type: 'admin'}])

这里groupBy和后面的mergeMap(toArray())都在switchMap的外层,它们的直接上游是switchMap返回的流,而switchMap的完成规则非常明确:

只有当它的上游(也就是示例中的arraySubject)完成,且当前正在运行的内部流也完成时,switchMap自身才会向下游发送完成信号。

示例中只调用了arraySubject.next()传入一次数组,始终没有调用arraySubject.complete(),因此switchMap会一直保持活跃,等待arraySubject的下一次输入,永远不会给下游的groupBy发送完成信号。
这种场景下,groupBy虽然把四个元素分成了user、admin两个分组,但两个分组Observable始终等不到完成信号,每个分组后挂载的toArray()就会一直缓存元素,永远不触发发射,自然订阅后收不到任何日志。
可以做个简单验证:在arraySubject.next(...)语句后加一句arraySubject.complete(),第一段代码就会正常输出分组结果。


第二段代码正常运行的原因

第二段代码把groupBy和mergeMap(toArray())放在了switchMap的内部投影函数里:

const arraySubject = new Subject();

arraySubject.pipe(
  switchMap((arr) => 
    from(arr).pipe(
      groupBy((elem) => elem.type),
      mergeMap((group) => group.pipe(toArray()))))
).subscribe((arrGroups) => {
  // 正常输出分组结果
})

arraySubject.next([{name: 'Bill', type: 'user'}, {name: 'Jane', type: 'admin'}, {name: 'John', type: 'user'}, {name: 'Ann', type: 'admin'}])

这里groupBy和toArray的直接上游是from(arr),而from(arr)是短生命周期流:把数组内的所有元素同步发送完成后,会立刻发送complete信号,执行逻辑如下:

  1. from(arr)把四个元素依次发给内部的groupBy后,立刻发送complete信号
  2. 内部的groupBy收到complete信号后,会把当前创建的user、admin两个分组Observable全部触发complete
  3. 每个分组后的toArray收到分组的complete信号,就会把缓存的同组元素打包成数组发射,因此可以正常拿到两个分组的结果

这个逻辑和外层arraySubject是否完成无关:每次arraySubject传入新的数组,switchMap都会创建一个全新的from(arr)流+独立的groupBy、toArray实例,处理完当前数组就自动完成,不会一直等待后续输入。


内容的提问来源于stack exchange,提问作者user2690235

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.03 10:45:39