RxJS并行及嵌套请求场景:登录请求执行成功但内部forEach循环未触发
问题原因
- RxJS 中所有 Observable 都是冷启动的,只有被调用
subscribe()方法、或者被其他上游 Observable 合并订阅时,才会执行内部的逻辑。你当前仅在forEach循环中创建了 Observable 实例,没有执行订阅操作,所以内部的链式逻辑完全不会运行。 - 混用 Promise 和 RxJS 操作符也会导致时序管控、异常处理的逻辑混乱,更推荐统一用 RxJS 操作符实现全链路逻辑。
修复方案
最小改动方案
仅需在每个 Observable 链式定义的末尾添加subscribe()即可触发执行:
entries.forEach((entry) => { return iif( // 原有iif逻辑不变 () => entry.extension == 'pdf', defer(() => this.convertFile.getPdfContentInText(entry.id)), defer(() => this.convertFile.getTextContent(entry.id)) ).pipe( // 原有pipe内的concatMap、catchError逻辑不变 concatMap((textContent: string) => { console.log(textContent); ctx.patchState({ qnaloading: true, }); return this.convertFile.convertFileContent(textContent); }), concatMap((convertedObject) => { return this.convertFile.uploadToQNA( entry.id, convertedObject, access_token ); }), concatMap((response: any) => { ctx.patchState({ loading: false }); return ctx.dispatch(new UpdateEntries([response.fileEntry])); }), catchError((e) => { ctx.patchState({ loading: false }); return of('reject: ' + JSON.stringify(e)); }) ).subscribe() // 新增这一行即可触发执行 });
更规范的全RxJS实现(推荐)
避免混用Promise,用RxJS操作符统一管控登录、并行请求逻辑,还能统一监听所有请求完成状态:
this.questionAnswerService.loginToQuriousApi().pipe( switchMap((response) => { const access_token = response.access_token // 把所有entry对应的Observable组装成数组 const entryTasks = entries.map((entry) => { return iif( () => entry.extension == 'pdf', defer(() => this.convertFile.getPdfContentInText(entry.id)), defer(() => this.convertFile.getTextContent(entry.id)) ).pipe( concatMap((textContent: string) => { console.log(textContent); ctx.patchState({ qnaloading: true }); return this.convertFile.convertFileContent(textContent); }), concatMap((convertedObject) => { return this.convertFile.uploadToQNA(entry.id, convertedObject, access_token); }), concatMap((response: any) => { return ctx.dispatch(new UpdateEntries([response.fileEntry])); }), catchError((e) => { return of('reject: ' + JSON.stringify(e)); }) ) }) // forkJoin并行执行所有entry的任务,所有任务完成后统一返回结果 return forkJoin(entryTasks) }), // 所有请求完成后统一更新状态 finalize(() => { ctx.patchState({ loading: false, qnaloading: false }); }) ).subscribe({ next: (results) => { // 所有任务执行完成后的回调,results包含所有任务的返回结果 console.log('所有文件处理完成', results) }, error: (e) => { console.log('全局异常', 'reject: ' + JSON.stringify(e)) } })
定位方法
你可以在iif的条件判断函数中插入console.log测试,如果日志没有打印,即可确认是 Observable 没有被订阅导致的逻辑未执行。
内容的提问来源于stack exchange,提问作者Sourabh Shah
相关产品推荐
相关产品推荐

