RxJS如何实现订阅追踪observable时才触发对应API observable请求
实现方案
该需求完全可以实现,RxJS的懒执行特性原生支持这种「订阅时才触发上游逻辑」的需求,核心思路是将API请求的触发时机从「方法调用时」延后到「返回的Observable被订阅时」,以下是几种可行的实现方案:
方案1:自定义Observable实现(最贴合预期逻辑)
直接将原有逻辑放到Observable构造函数的订阅回调中,只有当用户订阅返回的流时,才会执行内部的API请求逻辑,完全符合你期望的效果:
generateUrl$(uploadConfig: GenerateUploadUrlVariables, file: File) { // 返回自定义Observable,仅订阅时执行内部逻辑 return new Observable<ActionWithResult<{ url: string }>>((subscriber) => { const generation$ = new BehaviorSubject<ActionWithResult<{ url: string }>>({ state: 'INIT', result: null, }); // 将状态Subject的事件转发给外层订阅者 const generationSub = generation$.subscribe(subscriber); const complete = (result: ActionWithResult<{ url: string }>) => { generation$.next(result); generation$.complete(); }; // 仅订阅时才会触发此处的API请求 const requestSub = this.generateUploadUrl(uploadConfig).pipe( switchMap((result) => { generation$.next({ state: 'UPLOADING', result: null }); const url = result.data.generateUploadUrl.url || ''; return this.httpClient.pipe( tap(() => complete({ state: 'SUCCEEDED', result: { url } })), catchError((e) => { complete({ state: 'FAILED', result: null }); return throwError(() => e); }) ); }), catchError((e) => { complete({ state: ActionState.FAILED, result: null }); return throwError(() => e); }) ).subscribe(); // 退订时自动清理内部订阅,避免内存泄漏 return () => { generationSub.unsubscribe(); requestSub.unsubscribe(); }; }); }
调用generateUrl$但未订阅的场景下,内部API请求代码完全不会执行,不会产生资源浪费。
方案2:使用defer操作符简化实现
RxJS提供的defer操作符可以更简洁的实现懒执行效果,不需要手动管理订阅转发:
generateUrl$(uploadConfig: GenerateUploadUrlVariables, file: File) { return defer(() => { const generation$ = new BehaviorSubject<ActionWithResult<{ url: string }>>({ state: 'INIT', result: null, }); const complete = (result: ActionWithResult<{ url: string }>) => { generation$.next(result); generation$.complete(); }; this.generateUploadUrl(uploadConfig).pipe( switchMap((result) => { generation$.next({ state: 'UPLOADING', result: null }); const url = result.data.generateUploadUrl.url || ''; return this.httpClient.pipe( tap(() => complete({ state: 'SUCCEEDED', result: { url } })), catchError((e) => { complete({ state: 'FAILED', result: null }); return throwError(() => e); }) ); }), catchError((e) => { complete({ state: ActionState.FAILED, result: null }); return throwError(() => e); }) ).subscribe(); return generation$; }); }
defer会等到有订阅者订阅时,才执行内部函数生成实际的Observable流,效果和方案1完全一致。
优化方案:移除手动Subject维护
你还可以直接通过流操作符实现状态流转,省去手动维护Subject和订阅的步骤,代码更简洁,且天然支持懒执行:
generateUrl$(uploadConfig: GenerateUploadUrlVariables, file: File) { return this.generateUploadUrl(uploadConfig).pipe( switchMap(result => { const url = result.data.generateUploadUrl.url || ''; return this.httpClient.pipe( map(() => ({ state: 'SUCCEEDED', result: { url } })), startWith({ state: 'UPLOADING', result: null }), catchError(() => of({ state: 'FAILED', result: null })) ) }), startWith({ state: 'INIT', result: null }), catchError(() => of({ state: ActionState.FAILED, result: null })) ); }
内容的提问来源于stack exchange,提问作者Crocsx
相关产品推荐
相关产品推荐

