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(); } }); });
错误原因
forkJoin传入了无效值:track.assets.map的回调里没有返回Observable,默认返回undefined,而forkJoin要求传入Observable数组,这直接导致了报错。- 回调地狱(嵌套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
相关产品推荐
相关产品推荐

