使用windowCount+concatMap分批请求时仅处理首个窗口问题求助
问题原因与解决方案
核心原因
你的代码里range(1,10)是同步一次性发出所有10个元素的,而windowCount是通过热Observable(Subject)实现窗口分组的。当源Observable同步完成所有元素发射后,windowCount创建的所有窗口会立即完成。
而concatMap是依次订阅每个窗口Observable的:它先订阅第一个窗口,能拿到该窗口的元素并处理ajax请求;但当它处理完第一个窗口的所有请求、准备订阅第二个窗口时,第二个窗口已经处于完成状态(热Observable不会保留已发出的元素),所以无法获取到4、5、6这些元素,后续窗口同理,最终只得到第一个窗口的请求结果。
解决方案
有两种实用的解决方式:
1. 用bufferCount替代windowCount
bufferCount直接发出包含分组元素的数组,这些数组属于冷Observable的一部分,concatMap可以依次处理每个数组,不会出现元素丢失的问题:
import { range, bufferCount, map, concatMap, mergeMap } from 'rxjs'; import { ajax } from 'rxjs/ajax'; const ids = range(1, 10); const result = ids.pipe( bufferCount(3), // 替换windowCount为bufferCount concatMap((batch) => batch.pipe( mergeMap((id) => ajax(`https://jsonplaceholder.typicode.com/posts/${id}`).pipe( map((res) => res.response) ) ) ) ) ); result.subscribe((x) => console.log(x));
2. 保留windowCount,将窗口转为冷Observable
如果一定要用windowCount,可以通过toArray()把窗口内的元素转为数组,将热Observable转为冷Observable,确保后续订阅能获取到元素:
import { range, windowCount, map, concatMap, mergeMap, toArray } from 'rxjs'; import { ajax } from 'rxjs/ajax'; const ids = range(1, 10); const result = ids.pipe( windowCount(3), concatMap((win) => win.pipe( toArray(), // 将窗口元素转成数组,转为冷Observable mergeMap((batch) => batch.pipe( mergeMap((id) => ajax(`https://jsonplaceholder.typicode.com/posts/${id}`).pipe( map((res) => res.response) ) ) ) ) ) ) ); result.subscribe((x) => console.log(x));
进阶:批量请求完成后再处理下一批
如果希望每批请求全部完成后再处理下一批,并且获取每批的结果数组,可以结合forkJoin使用:
import { range, bufferCount, concatMap, forkJoin, map } from 'rxjs'; import { ajax } from 'rxjs/ajax'; const ids = range(1, 10); const result = ids.pipe( bufferCount(3), concatMap((batch) => forkJoin( batch.map(id => ajax(`https://jsonplaceholder.typicode.com/posts/${id}`).pipe( map(res => res.response) ) ) ) ) ); result.subscribe((batchResult) => console.log('批量结果:', batchResult));
内容的提问来源于stack exchange,提问作者Thore
相关产品推荐
相关产品推荐

