RxJS中mergeMap置于switchMap内与后置的执行差异问题
两者行为差异的核心是流的生命周期边界不同,本质是对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信号,执行逻辑如下:
from(arr)把四个元素依次发给内部的groupBy后,立刻发送complete信号- 内部的groupBy收到complete信号后,会把当前创建的user、admin两个分组Observable全部触发complete
- 每个分组后的toArray收到分组的complete信号,就会把缓存的同组元素打包成数组发射,因此可以正常拿到两个分组的结果
这个逻辑和外层arraySubject是否完成无关:每次arraySubject传入新的数组,switchMap都会创建一个全新的from(arr)流+独立的groupBy、toArray实例,处理完当前数组就自动完成,不会一直等待后续输入。
内容的提问来源于stack exchange,提问作者user2690235

