You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 10:01:18