如何正确对Observable列表使用forkJoin并处理订阅与完成回调?
并行Observable处理:实时更新结果+完成回调实现
问题背景
我有大量需要并行执行的有序Observable,每个Observable返回结果后,需将结果追加到BehaviorSubject中传递,但必须在所有Observable执行完成后调用特定函数。
这些Observable各自从API下载图片及相关元数据,要求尽可能快地执行请求,且每个结果一返回就立即处理,所有请求完成时触发完成回调——这意味着Observable需要并行执行。
原实现(无完成回调)
const requests: Observable[] = getRequests(); requests.forEach(obs => obs.subscribe(res => { const currentImages = this.$images.value; currentImages.push(res); this.$images.next(currentImages); }));
初次尝试的方案
为实现完成时的回调,我编写了以下代码,虽然功能可行,但将forkJoin与请求订阅拆分的写法显得冗余怪异,希望找到更优雅的实现方式:
const requests: Observable[] = getRequests(); const finishedTracker = new Subject<void>(); requests.forEach(obs => obs.subscribe(res => { const currentImages = this.$images.value; currentImages.push(res); this.$images.next(currentImages); })); forkJoin(requests).subscribe(() => { finishedTracker.next(); finishedTracker.complete(); console.log('requests done'); });
修正重复请求的尝试
后来意识到两次订阅会导致请求被执行两次,于是调整了实现方式:
from(requests).pipe( mergeMap(o => { o.subscribe(res => { const currentImages = this.$images.value; currentImages.push(res); this.$images.next(currentImages); }) return o; }, 10) ).subscribe(() => { finishedTracker.next(); console.log('requests done'); })
注:我没有直接使用forkJoin的结果,是因为它会等待所有请求完成后才返回全部结果,而我需要每个请求完成后立即将结果传递给BehaviorSubject——毕竟请求数量经常多达数百个。
最终采用的解决方案
from(requests).pipe( mergeMap(request => request, 10), scan<ImageResponse, ImageResponse[]>((all, current, index) => { all = all.concat(current); this.$images.next(all); return all; }, []) ).subscribe({ complete: () => { finishedTracker.next(); console.log('requests done'); } });
内容的提问来源于stack exchange,提问作者tgm
相关产品推荐
相关产品推荐

