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

RxJS使用带concurrent的mergeMap时Observable无法完成的原因

问题原因分析

这个异常行为的核心是异步timer任务与同步EMPTY处理在事件循环中的执行顺序冲突,导致RxJS的mergeMap无法正确跟踪所有内部Observable的完成状态:

  1. 源Observable from(Array(1000).keys()) 会同步发出0到999的所有值,mergeMap将这些值映射为100个timer(10)(异步)和900个EMPTY(同步完成),并把前10个timer(10)之外的所有内部Observable加入等待队列。
  2. 前10个timer(10)会在10ms后依次完成,mergeMap随后从队列中订阅下一批timer(10),这个过程持续到n=90对应的timer(10)完成。
  3. 当n=90的timer(10)完成后,mergeMap开始处理队列中的EMPTY。由于EMPTY是立即完成的,mergeMap会以同步方式快速处理剩下的900个EMPTY,这会占用事件循环的主线程。
  4. 而n=91到99的timer(10)回调属于宏任务,需要等待当前同步任务(处理所有EMPTY)完成后才能执行。在部分环境或RxJS旧版本中,大量同步的EMPTY处理会导致事件循环调度异常,这些宏任务无法被正常执行,mergeMap因此一直等待未完成的timer(10),最终无法触发complete回调。

至于你提到的几种正常情况:

  • 不指定concurrent时,mergeMap会同时订阅所有100个timer(10),它们的回调会批量执行,随后快速处理EMPTY,不会阻塞事件循环;
  • 移除条件或全返回EMPTY时,内部Observable要么全是能正常完成的异步任务,要么全是同步完成的任务,mergeMap能正确跟踪所有任务的完成状态。

解决方法

  1. 升级RxJS版本:这个问题在RxJS的新版本中已被修复,升级后即可正常运行。
  2. 给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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 03:54:59