如何将大量Observable网络请求按小批次执行?
解决Observable批量请求分批执行的问题
我有一个Observable数组,每个Observable对应一个API请求:
Observable<any>[] myWebRequests;
目前用combineLatest同时执行所有请求:
combineLatest(myWebRequests).subscribe(result => this.messageService.add({ severity: "success", summary: this.translateService.instant("StoryVersionEdit.PageTitle"), detail: this.translateService.instant("StoryVersionEdit.StoryVersionSaved"), closable: false }), error => { console.error(error); });
但当请求数量较多(比如80个)时,服务器会被大量并发请求压垮,导致部分请求失败。需要改成分批执行,每次处理10个请求,完成当前批次后再执行下一批。
先修正代码中的类型错误
首先注意到updateAllStoryVersionsBy方法的返回类型是Observable<StoryVersion>[],但实际代码里你把allStoryVersionsToUpdate.pipe(...)的结果强制转换成了数组——这是错误的,pipe操作后返回的是单个Observable(内部用forkJoin把多个请求合并成一个Observable),不是Observable数组。这个错误会导致后续的combineLatest使用异常,先修正这个问题。
分批处理的实现方案
核心思路:
- 将所有请求分成若干组,每组10个
- 用
concatMap依次处理每组,每组内部用forkJoin并行执行该组请求 - 确保上一批请求全部完成后,再启动下一批
步骤1:修改updateAllStoryVersionsBy返回正确的Observable
把方法返回类型改成Observable<StoryVersion[][]>,并去掉错误的类型强制转换:
updateAllStoryVersionsBy( assetHoldingNumber: string, assetHoldingLocation: string, originalValue_LegacyDataFormatSupport: string, originalValue_AssetHoldingEstablishment: string ): Observable<StoryVersion[][]> { // 保留原有验证逻辑 if (!this.storyVersion.assetHoldingNumber && !this.storyVersion.assetHoldingLocation) return of([]); const changedProperties = this.getChangedProperties(); // 保留原有属性检查逻辑 if (changedProperties.indexOf('assetHoldingNumber') < 0 && changedProperties.indexOf('assetHoldingLocation') < 0 && changedProperties.indexOf('assetHoldingEstablishment') < 0 && changedProperties.indexOf('assetHoldingRoom') < 0 && changedProperties.indexOf('assetHoldingNotes') < 0 && changedProperties.indexOf('recordingType') < 0 && changedProperties.indexOf('legacyItemHoldingType') < 0 && changedProperties.indexOf('legacyDataFormatSupport') < 0 && changedProperties.indexOf('cbcVideoResolution') < 0 ) return of([]); // 保留获取allStoryVersionsToUpdate的逻辑 var allStoryVersionsToUpdate: Observable<IGuid[]>; if (assetHoldingNumber) { allStoryVersionsToUpdate = this.storyVersionService.findStoryVersionsByAssetHoldingNumber(assetHoldingNumber); } else if (assetHoldingLocation) { var queryParameters: StoryVersionGetByAssetHoldingRequest = { isUniqueOnly: false, assetHoldingLocation: assetHoldingLocation, assetHoldingEstablishment: originalValue_AssetHoldingEstablishment, legacyItemHoldingType: originalValue_LegacyDataFormatSupport, assetHoldingNumber: null, responsible: null, select: null, offset: 0, limit: 100 }; allStoryVersionsToUpdate = this.storyVersionService.findStoryVersionsByAssetHoldingLocation(queryParameters); } return allStoryVersionsToUpdate.pipe( concatMap(data => { // 生成所有更新请求的Observable数组 const allRequests$ = data.map(g => { var svToUpdate: any = { guid: g.guid }; // 保留原有属性赋值逻辑 if (changedProperties.indexOf('assetHoldingNumber') >= 0) svToUpdate.assetHoldingNumber = this.storyVersionEditForm.value.assetHoldingNumber; if (changedProperties.indexOf('assetHoldingLocation') >= 0) svToUpdate.assetHoldingLocation = this.storyVersionEditForm.value.assetHoldingLocation; if (changedProperties.indexOf('assetHoldingEstablishment') >= 0) svToUpdate.assetHoldingEstablishment = this.storyVersionEditForm.value.assetHoldingEstablishment; if (changedProperties.indexOf('assetHoldingRoom') >= 0) svToUpdate.assetHoldingRoom = this.storyVersionEditForm.value.assetHoldingRoom; if (changedProperties.indexOf('assetHoldingNotes') >= 0) svToUpdate.assetHoldingNotes = this.storyVersionEditForm.value.assetHoldingNotes; if (changedProperties.indexOf('recordingType') >= 0) svToUpdate.recordingType = this.storyVersionEditForm.value.recordingType; if (changedProperties.indexOf('legacyItemHoldingType') >= 0) svToUpdate.legacyItemHoldingType = this.storyVersionEditForm.value.legacyItemHoldingType; if (changedProperties.indexOf('legacyDataFormatSupport') >= 0) svToUpdate.legacyDataFormatSupport = this.storyVersionEditForm.value.legacyDataFormatSupport; if (changedProperties.indexOf('cbcVideoResolution') >= 0) svToUpdate.cbcVideoResolution = this.storyVersionEditForm.value.cbcVideoResolution; return this.storyVersionService.updateStoryVersion(svToUpdate); }); // 将请求数组分成每10个一组 const batchSize = 10; const batches: Observable<StoryVersion[]>[] = []; for (let i = 0; i < allRequests$.length; i += batchSize) { const batch = allRequests$.slice(i, i + batchSize); batches.push(forkJoin(batch)); // 每组用forkJoin并行执行 } // 用concatMap依次执行所有批次,保证顺序 return from(batches).pipe( concatMap(batch$ => batch$), toArray() // 收集所有批次的结果 ); }) ); }
步骤2:调用并订阅处理结果
现在updateAllStoryVersionsBy返回的是单个Observable,直接订阅即可:
this.updateAllStoryVersionsBy(null, 'test', 'test', 'test').subscribe( allResults => { // allResults是二维数组,每个元素对应一个批次的请求结果 this.messageService.add({ severity: "success", summary: this.translateService.instant("StoryVersionEdit.PageTitle"), detail: this.translateService.instant("StoryVersionEdit.StoryVersionSaved"), closable: false }); }, error => { console.error(error); // 可在此添加错误提示,比如批次失败后的处理 } );
关键操作说明
from(batches):把批次Observable数组转换成流,每次发射一个批次的ObservableconcatMap(batch$ => batch$):依次订阅每个批次,只有当前批次全部完成才会启动下一批forkJoin(batch):并行执行当前批次内的所有请求,需所有请求成功才返回结果(若要容错,可给单个请求添加catchError)
如果需要允许单个请求失败不影响整个批次,可修改单个请求的错误处理:
return this.storyVersionService.updateStoryVersion(svToUpdate).pipe( catchError(err => { console.error(`更新失败: ${g.guid}`, err); return of(null); // 返回null标记失败,后续可过滤掉 }) );
内容的提问来源于stack exchange,提问作者Bruno
相关产品推荐
相关产品推荐

