如何在forEach循环中多次触发Subject?Angular RxJS问题求助
核心问题:switchMap的特性导致请求被取消
你遇到的问题根源在于**switchMap的行为特性**:每当源Observable(这里是selectedVersion$)发出新值时,switchMap会立即取消之前正在进行的内部Observable(也就是this.leadTimeService.getLeadTimeVersionData(...)的请求),转而订阅新的内部Observable。
你的this.versions.forEach(v => this.selectedVersion$.next(v))是同步快速触发7次next,此时前几次的请求还没来得及返回结果,就被后面的next触发的新请求给"切换"掉了,最终只有最后一次请求的结果被订阅接收到,所以counter最终只会等于1。
至于你尝试的setTimeout无效,是因为即使把next放到异步执行,只要请求的响应速度跟不上next的触发速度,switchMap依然会取消未完成的前序请求,除非你把延迟设置得足够长(比如超过每个请求的响应时间),但这显然不是合理的解决方案。
针对不同需求的解决方案
方案1:并行发起所有请求,等待全部完成后合并结果(推荐,性能最优)
如果你不需要按顺序处理请求,只是想拿到所有版本的数据合并到this.nodes,可以用forkJoin一次性并行发起所有请求,等全部完成后统一处理:
this.leadTimeService.selectedLeadTimeWorkbench$.pipe( takeUntil(this.unsubscribe), switchMap(x => { this.productLineId = x.id; this.versions = x.workingVersions; // 把每个版本转换成对应的请求Observable const versionRequests = this.versions.map(v => this.leadTimeService.getLeadTimeVersionData(this.productLineId, v.id) ); // 并行发起所有请求,全部完成后返回结果数组 return forkJoin(versionRequests); }) ).subscribe(results => { // 遍历所有结果,合并到nodes results.forEach(res => { this.nodes = mergeArray(this.nodes, res.nodes); }); // 所有请求完成后执行过滤 this.filterTables(''); });
方案2:按顺序依次发起请求(前一个完成再发下一个)
如果需要保证请求的顺序(比如依赖前一个请求的结果),可以把switchMap换成concatMap,它会等待前一个内部Observable完成后,再订阅下一个:
// 修改selectedVersion$的订阅部分 this.selectedVersion$.pipe( concatMap(v => this.leadTimeService.getLeadTimeVersionData(this.productLineId, v.id)), // 替换switchMap为concatMap takeUntil(this.unsubscribe) ).subscribe(x => { this.nodes = mergeArray(this.nodes, x.nodes); counter++; if(counter === this.versions.length){ this.filterTables(''); } }); // 触发部分保持不变 this.leadTimeService.selectedLeadTimeWorkbench$.pipe( takeUntil(this.unsubscribe) ).subscribe(x => { this.productLineId = x.id; this.versions = x.workingVersions; this.versions.forEach(v => { this.selectedVersion$.next(v); }); });
方案3:保留并行处理,但不取消前序请求(使用mergeMap)
如果你希望同时处理多个请求,且不取消前序请求(所有请求的结果都会被接收),可以把switchMap换成mergeMap,还可以通过第二个参数控制并发数(比如限制同时发起3个请求):
this.selectedVersion$.pipe( mergeMap(v => this.leadTimeService.getLeadTimeVersionData(this.productLineId, v.id), 3), // 第二个参数是并发数 takeUntil(this.unsubscribe) ).subscribe(x => { this.nodes = mergeArray(this.nodes, x.nodes); counter++; if(counter === this.versions.length){ this.filterTables(''); } }); // 触发部分保持不变
补充说明
switchMap适合需要取消前序请求的场景(比如搜索框输入,只保留最新的搜索请求);concatMap适合按顺序执行的场景;mergeMap适合并行执行且保留所有结果的场景;forkJoin适合一次性发起多个请求,等待全部完成的场景,是并行处理的最优选择之一。
内容的提问来源于stack exchange,提问作者Allen

