RxJS动态添加/移除Observable:可插拔ProgressService实现问询
实现可插拔的响应式ProgressService
针对你想要构建的动态可插拔ProgressService,我们可以利用RxJS的Subject和操作符来实现完全响应式的动态流管理,不需要手动调用increase()/decrease(),同时支持动态添加/移除加载提供者,还能处理你提到的优化细节。
核心实现思路
我们用一个Subject来维护所有注册的加载提供者Observable,通过mergeAll()动态合并这些流,再用scan计算当前加载任务的总数。同时为每个提供者添加清理逻辑,确保移除时能正确更新计数。
完整代码实现
import { Injectable } from '@angular/core'; import { Observable, Subject, of, scan, map, mergeAll, takeUntil, finalize, startWith, distinctUntilChanged, filter } from 'rxjs'; @Injectable({ providedIn: 'root' }) export class ProgressService { // 用于接收动态添加的加载提供者流 private providers$ = new Subject<Observable<boolean>>(); // 维护当前加载任务总数的流 private progress$: Observable<number>; constructor() { this.progress$ = this.providers$.pipe( mergeAll(), // 累加/递减加载任务数,确保不会出现负数 scan((total, isLoading) => { const newValue = total + (isLoading ? 1 : -1); return Math.max(newValue, 0); }, 0), // 避免重复发射相同数值,减少不必要的UI更新 distinctUntilChanged() ); } /** * 添加加载提供者 * @param lp 发射true表示开始加载,false表示加载结束的Observable * @returns 注销该提供者的函数,调用后停止监听并清理计数 */ addLoadingProvider(lp: Observable<boolean>): () => void { const cancel$ = new Subject<void>(); let lastState = false; // 处理加载提供者流:忽略首次false,记录状态,添加清理逻辑 const processedProvider$ = lp.pipe( startWith(null), distinctUntilChanged(), // 过滤首次发射的false:只有之前是加载状态,或者当前是加载状态时才发射 filter(val => val !== null && (lastState || val)), map(val => { lastState = val; return val; }), // 当流结束(包括主动注销)时,如果最后状态是加载中,自动递减计数 finalize(() => { if (lastState) { this.providers$.next(of(false)); } }), takeUntil(cancel$) ); this.providers$.next(processedProvider$); // 返回注销函数,供外部调用 return () => { cancel$.next(); cancel$.complete(); }; } /** * 返回是否正在加载的Observable */ isLoading(): Observable<boolean> { return this.progress$.pipe( map(total => total > 0) ); } }
关键细节说明
动态流管理:
- 用
providers$这个Subject来接收新的加载提供者,mergeAll()会自动合并所有传入的流,实现动态添加的效果。 - 不需要单独的
removeLoadingProvider方法,而是让addLoadingProvider返回一个注销函数,符合RxJS的订阅管理习惯,用户可以自行保存并在需要时调用。
- 用
首次false过滤:
通过startWith(null)+filter逻辑,确保如果一个加载提供者首次发射的是false,不会触发计数递减,避免无效的状态变更。注销时的清理:
每个提供者流都添加了finalize钩子,当流结束(包括主动注销)时,如果最后状态是加载中(lastState = true),会自动向providers$发射一个false,确保计数正确递减。计数安全:
在scan中用Math.max(newValue, 0)确保计数不会变成负数,避免因异常情况导致的状态错误。
使用示例
// 在组件中使用 export class DataLoadingComponent { private cancelLoading?: () => void; isLoading$ = this.progressService.isLoading(); constructor( private progressService: ProgressService, private dataService: DataService ) {} startFetchingData() { // 构造加载状态流:开始时发射true,请求完成后发射false const loadState$ = this.dataService.fetchData().pipe( map(() => false), startWith(true) ); // 添加加载提供者并保存注销函数 this.cancelLoading = this.progressService.addLoadingProvider(loadState$); } abortFetch() { // 主动注销加载提供者 this.cancelLoading?.(); } ngOnDestroy() { // 组件销毁时清理 this.abortFetch(); } }
这个实现完全基于响应式思想,避免了手动维护状态的繁琐和潜在错误,同时完美支持动态添加/移除加载提供者的需求。
内容的提问来源于stack exchange,提问作者Benjamin M
相关产品推荐
相关产品推荐

