【问题标题】:How to create an Observable which updates it's value every 10 seconds (interval) only if subscribed?如何创建一个仅在订阅时每 10 秒(间隔)更新其值的 Observable?
【发布时间】:2019-08-22 15:07:03
【问题描述】:

我正在使用 RxJS 构建一个 Angular 应用程序。在我的 StoreService 中,我想要一个 Observable,我的组件可以从中订阅。可观察的 public unreadMessagesCount: Observable<number> 保存来自 api 的未读消息的值。

应该每 10 秒完成一次 api 调用,并且 observable 应该使用新值更新自身。但是,如果没有组件订阅此 observable,则不应进行轮询。当有人订阅时,应该立即完成。这样数据从一开始就是最新的。

有没有一种优雅的方式来使用 RxJS 构建它?提前致谢!

【问题讨论】:

    标签: angular rxjs ngrx


    【解决方案1】:

    您可以使用intervalswitchMap 创建一个可观察的对象,当您订阅它时将开始轮询您的 API:

    import { interval } from 'rxjs';
    
    const unreadMessagesCount$ = interval(10000).pipe(switchMap(_ => apiCall()));
    

    但是,对于您执行的每个订阅,API 调用都会重复。如果您不希望这种情况发生,则必须在所有订阅者之间共享可观察的源。您可以为此使用multicast

    import { interval, Subject } from 'rxjs';
    import { multicast} from 'rxjs/operators';
    
    const unreadMessagesCount$ = interval(10000).pipe(
      switchMap(_ => apiCall()),
      multicast(() => new Subject())
    );
    

    那么你只需要在 observable 上调用connect

    unreadMessagesCount$.subscribe(value => ...);
    unreadMessagesCount$.subscribe(value => ...);
    
    unreadMessagesCount$.connect();
    

    this Stackblitz demo

    【讨论】:

      【解决方案2】:

      使用timer 每 N 秒创建一个可观察的发射,从 M 秒开始:

      interval$ = timer(1, 10000);
      

      然后您可以声明您的 API 调用:

      APICall$ = this.http.get('...')
      

      最后,从这两个中创建一个 observable:

      pulledData$ = this.interval$.pipe(
        switchMap(() => this.APICall$)
      );
      

      现在,除非您订阅它,否则您将没有任何价值。但是一旦你订阅了它,你就会进行一次 API 调用(尽管是在 1 毫秒之后),并且每十秒钟你就会更新一次调用。请记住在销毁组件时取消订阅!

      【讨论】:

        【解决方案3】:

        您可以使用timer 并定期使用concatMap 进行异步调用。如果您希望它在没有初始延迟的情况下发出,您可以使用timer(0, 10000) 语法:

        import { timer } from 'rxjs';
        import { concatMap } from 'rxjs/operators';
        
        timer(0, 10000).pipe(
          concatMap(i => makeHTTPCall(i)),
        );
        

        【讨论】:

        • 很简单,谢谢! makeHTTPCall() 应该返回什么?
        • makeHTTPCall 需要返回一个 Observable 或 Promise。
        【解决方案4】:

        您可以使用SubjectBehaviorSubjecttimer

        import { BehaviorSubject, Observable, timer } from 'rxjs';
        import { tap } from 'rxjs/operators';
        
        ...
        export class YourService {
          value$ = new BehaviorSubject('your initial value');
        
          constructor() {
            timer(10000).pipe(
              tap(() => this.loadValue())
            ).subscribe();
          }
        
          loadValue() {
           callYourAPI().subscribe(
             value => this.value$.next(value);
           );
          }
        }
        
        

        【讨论】:

        • 不需要行为主体并进行初始订阅,只需在服务中有一个公共可观察对象,其值为timer() 可观察对象,只要组件需要订阅,它们就会订阅它跨度>
        • 是的,你是对的,我的错。想了一种方法来手动触发重新加载
        猜你喜欢
        • 1970-01-01
        • 2019-05-29
        • 1970-01-01
        • 1970-01-01
        • 2015-10-21
        • 1970-01-01
        • 1970-01-01
        • 2017-10-18
        • 2020-12-06
        相关资源
        最近更新 更多