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

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)
    );
  }
}

关键细节说明

  1. 动态流管理:

    • 用providers$这个Subject来接收新的加载提供者,mergeAll()会自动合并所有传入的流,实现动态添加的效果。
    • 不需要单独的removeLoadingProvider方法,而是让addLoadingProvider返回一个注销函数,符合RxJS的订阅管理习惯,用户可以自行保存并在需要时调用。
  2. 首次false过滤:
    通过startWith(null) + filter逻辑,确保如果一个加载提供者首次发射的是false,不会触发计数递减,避免无效的状态变更。

  3. 注销时的清理:
    每个提供者流都添加了finalize钩子,当流结束(包括主动注销)时,如果最后状态是加载中(lastState = true),会自动向providers$发射一个false,确保计数正确递减。

  4. 计数安全:
    在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 08:11:51