RxJS使用带concurrent的mergeMap时Observable无法完成的原因
问题原因分析
这个异常行为的核心是异步timer任务与同步EMPTY处理在事件循环中的执行顺序冲突,导致RxJS的mergeMap无法正确跟踪所有内部Observable的完成状态:
- 源Observable
from(Array(1000).keys())会同步发出0到999的所有值,mergeMap将这些值映射为100个timer(10)(异步)和900个EMPTY(同步完成),并把前10个timer(10)之外的所有内部Observable加入等待队列。 - 前10个
timer(10)会在10ms后依次完成,mergeMap随后从队列中订阅下一批timer(10),这个过程持续到n=90对应的timer(10)完成。 - 当n=90的
timer(10)完成后,mergeMap开始处理队列中的EMPTY。由于EMPTY是立即完成的,mergeMap会以同步方式快速处理剩下的900个EMPTY,这会占用事件循环的主线程。 - 而n=91到99的
timer(10)回调属于宏任务,需要等待当前同步任务(处理所有EMPTY)完成后才能执行。在部分环境或RxJS旧版本中,大量同步的EMPTY处理会导致事件循环调度异常,这些宏任务无法被正常执行,mergeMap因此一直等待未完成的timer(10),最终无法触发complete回调。
至于你提到的几种正常情况:
- 不指定
concurrent时,mergeMap会同时订阅所有100个timer(10),它们的回调会批量执行,随后快速处理EMPTY,不会阻塞事件循环; - 移除条件或全返回
EMPTY时,内部Observable要么全是能正常完成的异步任务,要么全是同步完成的任务,mergeMap能正确跟踪所有任务的完成状态。
解决方法
- 升级RxJS版本:这个问题在RxJS的新版本中已被修复,升级后即可正常运行。
- 给EMPTY添加异步调度:让
EMPTY的完成通知进入宏任务队列,避免同步阻塞事件循环:
import { from, timer, mergeMap, EMPTY, delay } from "rxjs"; from(Array(1000).keys()).pipe( mergeMap(n => n < 100 ? timer(10) : EMPTY.pipe(delay(0)), 10) ).subscribe({ complete: () => console.log('Complete') });
内容的提问来源于stack exchange,提问作者Shigure
相关产品推荐
相关产品推荐

