RxJS如何将Observable<Observable<Thing>>转换为Observable<Thing[]>
问题原因
你当前的写法会返回非预期类型,核心是mergeMap处理Observable<Thing>[]类型的返回值时,会将数组内的可观察对象逐个打平订阅,最终流输出的是独立的Thing实例,而非聚合后的结果数组。
解决方案
直接在map生成的单个请求可观察对象数组外层包裹forkJoin操作符即可,修改后代码如下:
getThings = (): Observable<Thing[]> => this.getThingIds().pipe( mergeMap((thingIds: number[]) => forkJoin(thingIds.map((id: number) => this.http.get<Thing>(`url${id}`))) ) );
forkJoin会等待数组内所有请求完成,一次性发射包含所有请求结果的数组,外层mergeMap将forkJoin的输出打平后,整个流的返回类型就是你需要的Observable<Thing[]>。
可选优化方案
如果需要控制请求并发数,避免同一时间发起太多请求触发接口限流,可以用拆分写法调整并发:
getThings = (): Observable<Thing[]> => this.getThingIds().pipe( mergeMap((thingIds: number[]) => from(thingIds)), // 第二个参数设为并发数,比如2代表最多同时发起2个请求 mergeMap((id: number) => this.http.get<Thing>(`url${id}`), 2), toArray() );
注意事项
forkJoin和上述拆分写法默认都会在任意一个请求报错时中断整个流,如果你需要允许单个请求失败不影响整体结果,可以给单个请求增加错误处理:
this.http.get<Thing>(`url${id}`).pipe( catchError(err => { // 可自定义错误场景返回值,比如null或者包含错误信息的对象 return of(null) }) )
内容的提问来源于stack exchange,提问作者tbarbot
相关产品推荐
相关产品推荐

