Rxjs实现Observable数组按批次顺序执行并合并结果
用RxJS实现Observable数组的顺序批次执行并合并结果
核心思路
将Observable数组按指定大小拆分为多个批次,顺序执行每个批次(前一批次全部完成后才启动下一批次),每个批次内的请求并行处理,最终将所有批次的结果合并为一个数组一次性返回。
实现代码
1. 数组拆批工具函数
先实现一个简单的数组拆分函数,将原Observable数组分成指定大小的批次:
function chunkArray<T>(arr: T[], chunkSize: number): T[][] { const chunks = []; for (let i = 0; i < arr.length; i += chunkSize) { chunks.push(arr.slice(i, i + chunkSize)); } return chunks; }
2. RxJS核心处理逻辑
利用from、concatMap、forkJoin和reduce组合实现需求:
import { from, forkJoin, Observable } from 'rxjs'; import { concatMap, reduce } from 'rxjs/operators'; // 你的原Observable数组 const responses$: Observable<Response>[] = [ this.service.get(1), this.service.get(2), this.service.get(3), this.service.get(4) ]; const batchSize = 2; // 自定义批次大小 const batches = chunkArray(responses$, batchSize); // 生成最终的结果Observable const allResults$ = from(batches).pipe( // 顺序处理每个批次,确保前一批完成后再执行下一批 concatMap(batch => forkJoin(batch)), // 将所有批次的结果合并为一个数组 reduce((accumulator, batchResults) => [...accumulator, ...batchResults], [] as Response[]) ); // 订阅获取合并后的全部结果 allResults$.subscribe({ next: (allResponses) => { // 这里拿到所有请求的结果,顺序与原数组一致 console.log(allResponses); }, error: (error) => { console.error('请求执行失败:', error); } });
关键操作说明
chunkArray:手动拆分数组,无额外依赖,适用于任意类型数组的批处理场景。from(batches):将批次数组转换为Observable流,逐个输出每个批次的Observable数组。concatMap:强制Observable流按顺序执行,只有当前批次的forkJoin完成后,才会处理下一个批次,从根本上控制并发量。forkJoin:并行执行单个批次内的所有请求,当批次内所有Observable都完成时,返回该批次的结果数组。reduce:将每个批次的结果数组逐步合并为一个大数组,最终一次性返回所有请求的结果。
针对大数量请求的适配
对于2000+请求的场景,只需调整batchSize参数(比如设置为10或20),即可控制同时运行的请求数量,既保证处理效率,又避免超出API的并发限制。整个实现完全基于RxJS,无需额外依赖。
内容的提问来源于stack exchange,提问作者Nobita Gascón
相关产品推荐
相关产品推荐

