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

Angular新手用嵌套forkJoin替代Promise.all遇RxJS undefined错误求助

问题描述

我是Angular和RxJS新手,之前习惯用Promise,现在要实现类似Promise.all()的功能:现有一个数组,每个元素需要发起HTTP请求,且每个元素下还有子数组,子数组的每个元素也得发起请求。之前单独用forkJoin没问题,但嵌套使用后一直报错:You provided "undefined" where a stream was expected。所有API调用都返回Observable,怀疑是forkJoin的使用方式有问题,代码如下:

.updateTrackSystemCompact({ ...uploadTrackSystemCompact, trackIds: null })
.subscribe((uploadInspectionJob: UploadInspectionJob) => {
  uploadTrackSystemCompact = {
    id: uploadTrackSystem.id,
    country: uploadTrackSystem.country,
    standardsUsed: uploadTrackSystem.standardsUsed,
    region: uploadTrackSystem.region,
    company: uploadTrackSystem.company,
    systemId: uploadTrackSystem.systemId,
    postalCode: uploadTrackSystem.postalCode,
    address: uploadTrackSystem.address,
    city: uploadTrackSystem.city,
    state: uploadTrackSystem.state,
    createdAt: uploadTrackSystem.createdAt,
    updatedAt: uploadTrackSystem.updatedAt,
    inspectionReport: uploadTrackSystem.inspectionReport,
    inspectionReports: uploadTrackSystem.inspectionReports,
    trackIds: trackIndexes,
  };
  this.trackSystemService
    .doUploadS3TrackSystemCompact(uploadTrackSystemCompact, uploadInspectionJob.uploadUrl)
    .subscribe((resp) => {
      if (resp === undefined) {
        let alert = this.alertCtrl.create({
          title: "S3 response was undefined",
          message: `response status is ${resp.status}`,
          buttons: ["OK"],
        });
        alert.present();
      }
      if (resp !== undefined) {
        forkJoin(
          uploadTrackSystem.tracks.map((track, trackIndex) => {
            const inspectionJobId = uploadInspectionJob.id;
            forkJoin(
              track.assets.map((asset, assetIndex) => {
                const request: UploadCompactRequestAsset = {
                  inspectionJobId: inspectionJobId,
                  trackIndex: trackIndex,
                  assetIndex: assetIndex,
                };
                this.trackSystemService
                  .getAssetUploadUrl(request)
                  .subscribe((response: UploadCompactResponse) => {
                    const url = response.uploadUrl;
                    this.trackSystemService.uploadS3Asset(asset, url);
                  });
              })
            ).subscribe(() => {
              const request: UploadCompactRequestTrack = {
                inspectionJobId: inspectionJobId,
                trackIndex: trackIndex,
              };
              this.trackSystemService
                .getTrackUploadUrl(request)
                .subscribe((response: UploadCompactResponse) => {
                  const url = response.uploadUrl;

                  const uploadTrackCompact: TrackCompact = {
                    id: track.id,
                    uuid: track.uuid,
                    trackId: track.trackId,
                    inspectionId: track.inspectionId,
                    trackSystemId: track.trackSystemId,
                    start: track.start,
                    end: track.end,
                    trackReport: track.trackReport,
                    trackReports: track.trackReports,
                    assetIds: track.assets.map((asset, index) => {
                      return index;
                    }),
                  };
                  this.trackSystemService.uploadS3TrackCompact(uploadTrackCompact, url);
                });
            });
          })
        ).subscribe(() => {
          this.trackSystemService.setReadyStatus(uploadInspectionJob.id).subscribe(
            (inspectionJob: InspectionJob) => {
              this.dismissLoading();
              // this.uploadButtonDisabled = false;
              this.awaitInspectionProcessing();
            },
            (err) => {
              this.increaseFailedUploadCount();
              this.uploadButtonDisabled = false;
              this.dismissLoading();
              switch (err.status) {
                case 401:
                  Utils.sessionExpiredAlert(this.authService, this.alertCtrl, this.navCtrl);
                  break;

                default:
                  let alert = this.alertCtrl.create({
                    title: "Error",
                    subTitle: `There was an error setting the ready status: (${err.status})`,
                    buttons: ["OK"],
                  });
                  alert.present();
                  break;
              }
            }
          );
        });
      } else {
        this.increaseFailedUploadCount();
        this.uploadButtonDisabled = false;
        this.dismissLoading();
        let alert = this.alertCtrl.create({
          title: "Error",
          subTitle: "There was an error uploading to the S3",
          buttons: ["OK"],
        });
        alert.present();
      }
    });
});
错误原因
  1. forkJoin传入了无效值:track.assets.map的回调里没有返回Observable,默认返回undefined,而forkJoin要求传入Observable数组,这直接导致了报错。
  2. 回调地狱(嵌套subscribe):你用了多层嵌套的subscribe,这违背了RxJS的设计理念,不仅代码可读性差,还容易引发内存泄漏和错误处理混乱的问题。
解决方案

用RxJS的操作符组合流,替代嵌套subscribe,同时确保每个map回调都返回合法的Observable。核心逻辑是:先完成每个track下所有asset的上传,再完成该track的上传,最后统一处理所有track完成后的逻辑。

修正后的代码:

.updateTrackSystemCompact({ ...uploadTrackSystemCompact, trackIds: null })
.pipe(
  switchMap((uploadInspectionJob: UploadInspectionJob) => {
    // 更新uploadTrackSystemCompact对象
    uploadTrackSystemCompact = {
      ...uploadTrackSystem,
      trackIds: trackIndexes
    };
    // 先上传trackSystemCompact到S3
    return this.trackSystemService.doUploadS3TrackSystemCompact(uploadTrackSystemCompact, uploadInspectionJob.uploadUrl)
      .pipe(
        // 只有S3上传成功才继续
        filter(resp => resp !== undefined),
        // 触发所有track的上传流程
        switchMap(() => {
          const trackObservables = uploadTrackSystem.tracks.map((track, trackIndex) => {
            const inspectionJobId = uploadInspectionJob.id;
            // 先处理当前track下的所有asset上传
            const assetObservables = track.assets.map((asset, assetIndex) => {
              const request: UploadCompactRequestAsset = {
                inspectionJobId,
                trackIndex,
                assetIndex
              };
              return this.trackSystemService.getAssetUploadUrl(request).pipe(
                switchMap((response: UploadCompactResponse) => {
                  // 获取上传URL后,执行asset上传
                  return this.trackSystemService.uploadS3Asset(asset, response.uploadUrl);
                })
              );
            });

            // 等当前track的所有asset上传完成后,再上传track自身
            return forkJoin(assetObservables).pipe(
              switchMap(() => {
                const request: UploadCompactRequestTrack = {
                  inspectionJobId,
                  trackIndex
                };
                return this.trackSystemService.getTrackUploadUrl(request).pipe(
                  switchMap((response: UploadCompactResponse) => {
                    const uploadTrackCompact: TrackCompact = {
                      ...track,
                      assetIds: track.assets.map((_, index) => index)
                    };
                    return this.trackSystemService.uploadS3TrackCompact(uploadTrackCompact, response.uploadUrl);
                  })
                );
              })
            );
          });

          // 等待所有track的上传流程完成
          return forkJoin(trackObservables).pipe(
            // 最后设置ready状态
            switchMap(() => this.trackSystemService.setReadyStatus(uploadInspectionJob.id))
          );
        }),
        // 统一错误处理
        catchError(err => {
          this.increaseFailedUploadCount();
          this.uploadButtonDisabled = false;
          this.dismissLoading();
          
          if (err.status === 401) {
            Utils.sessionExpiredAlert(this.authService, this.alertCtrl, this.navCtrl);
          } else {
            const alert = this.alertCtrl.create({
              title: "Error",
              subTitle: `There was an error: (${err.status})`,
              buttons: ["OK"]
            });
            alert.present();
          }
          return throwError(() => err);
        })
      );
  })
)
.subscribe({
  next: (inspectionJob: InspectionJob) => {
    this.dismissLoading();
    this.awaitInspectionProcessing();
  },
  error: () => {
    // 错误已在catchError中处理,这里可以留空或补充额外逻辑
  }
});

关键改进点:

  • 用pipe和操作符(switchMap、forkJoin、filter、catchError)替代嵌套subscribe,代码更清晰
  • 每个map回调都返回Observable,确保forkJoin能正确接收流
  • 统一错误处理,避免重复的错误逻辑
  • 利用对象扩展符简化对象赋值,提升代码可读性

内容的提问来源于stack exchange,提问作者Ryan Speciale

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 16:35:04