调用next时Subject输出顺序异常:带副作用的map操作符场景
分析Subject同步递归next导致的顺序问题
嘿,我来帮你拆解这个问题的核心原因,以及怎么解决它~
首先,咱们先明确你遇到的问题本质:你在Subject的订阅回调(准确说是pipe里的tap操作符)中同步调用了sub.next(),而Subject的next()是同步执行所有订阅者的,这就导致了递归式的执行栈嵌套,最终打乱了你预期的输出顺序。
具体执行流程拆解
咱们拿你给出的代码为例,假设你一开始调用了sub.next(2),实际的执行顺序是这样的:
- 调用
sub.next(2),Subject会依次通知所有订阅者:- 先执行第一个订阅
emit$的回调:把2推入emits数组 →emits = [2] - 接着执行第二个订阅
data$的pipe流程:map返回2,第一个tap把2推入mapped→mapped = [2]- 进入第二个
tap,判断2是偶数,立刻调用sub.next(2333) - 这时候,当前的
2的pipe流程会被打断,立刻执行sub.next(2333)的所有订阅逻辑:- 先执行
emit$:emits.push(2333)→emits = [2, 2333] - 再执行
data$的pipe:map返回2333,第一个tap推入mapped→mapped = [2, 2333]- 第二个
tap判断x===2333,调用sub.next(1111) - 再次打断当前流程,执行
sub.next(1111)的订阅:emit$:emits.push(1111)→emits = [2, 2333, 1111]data$的pipe处理1111,没有触发新的next,流程完成
- 回到
2333的pipe流程,完成剩余步骤
- 先执行
- 回到
2的pipe流程,完成剩余步骤
- 先执行第一个订阅
你看,整个过程是嵌套递归执行的,而不是你可能预期的“先把2的所有处理做完,再处理2333,最后处理1111”的线性顺序。这就是为什么输出顺序不符合预期的根本原因。
解决方案:异步触发新的next
如果想要让新的next在当前值的所有处理流程完成后再执行,你需要把sub.next()的调用放到异步队列里,避免同步递归。RxJS提供了asyncScheduler来做这件事,或者也可以用原生的setTimeout(不过更推荐用RxJS的调度器保持一致性)。
修改后的代码示例:
import { Subject } from 'rxjs'; import { map, tap } from 'rxjs/operators'; import { asyncScheduler } from 'rxjs'; const sub = new Subject(); const emits = []; const mapped = []; const emit$ = sub.asObservable().subscribe(x => emits.push(x)); const data$ = sub.asObservable().pipe( map(x => { return x; }), tap(x => mapped.push(x)), tap(x => { if (x % 2 === 0) { // 用asyncScheduler异步触发next asyncScheduler.schedule(() => sub.next(2333)); } if (x === 2333) { asyncScheduler.schedule(() => sub.next(1111)); } }) ).subscribe(); // 初始触发 sub.next(2);
这样修改后,执行顺序就会变成线性的:
- 先完成
2的所有订阅处理(emits和mapped都推入2) - 当前执行栈清空后,执行异步队列里的
sub.next(2333),完成它的所有处理 - 再执行异步队列里的
sub.next(1111),完成它的所有处理
最终emits和mapped的顺序都会是[2, 2333, 1111],符合常规的预期。
关键知识点回顾
- Subject的
next()是同步执行所有订阅者回调的,这意味着在回调中调用next()会立刻中断当前流程,优先处理新的next事件 - 同步递归调用
next会导致执行栈嵌套,打乱预期的顺序 - 若要保证顺序,需将新的
next调用延迟到当前执行栈之外,使用异步调度器是RxJS中的标准做法
内容的提问来源于stack exchange,提问作者Tao Zhu
相关产品推荐
相关产品推荐

