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

如何将大量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数组转换成流,每次发射一个批次的Observable
  • concatMap(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 09:35:17