【问题标题】:How to use one RXJs stream for 2 different things如何将一个 RXJ 流用于两件不同的事情
【发布时间】:2016-02-12 16:49:15
【问题描述】:

我的应用程序中有 2 个服务。一个通过网络加载 YouTube 评论线程和 cmets 的 YouTube 服务,以及一个管理批量加载 cmets 的评论服务。 cmets 服务特定于我的应用程序,youtube 服务与应用程序无关。

我有一个函数getCommentThreadsForChannel 可以加载评论线程。主要的实现是在 youTube 服务中,但这是在评论服务上调用的,它基本上只是调用返回来自 youTube 服务的 observable。

就我调用它的控制器而言,这只是一个可观察的评论线程序列。但是,在我的 commentService 中,我想将这些线程存储在本地。每当我获得 100 个以上的线程时,我想将其批量存储到所有线程中,这样我就不会处理每个新数据位的列表。我想出了这个代码:

    getCommentThreadsForChannel(): Rx.Observable<ICommentThread> {
        var threadStream: Rx.Observable<ICommentThread> =
            this.youTubeService.getCommentThreadsForChannel();

        threadStream
            .bufferWithCount(100)
            .scan( ( allItems, currentItem ) => {
                currentItem.forEach(thread => {
                    allItems.push(thread);
                });

                console.log( `Save items to local storage: ${allItems.length}` )

                return allItems;
            }, []  );

        return threadStream;
    }

我认为这里用于批处理线程并将所有线程累积到一个数组中的逻辑很好,但从未调用此代码。我想这是因为我根本没有订阅这个帖子。

我不想在这里订阅,因为这将订阅底层流,然后我将有 2 个订阅,所有数据将加载两次(有很多数据 - 加载所有数据大约需要一分钟一次超过 30 次调用 100 个线程)。

这里我基本上想要一个do,它不会影响传递给控制器​​的流,但我想使用RXjs的缓冲和累积逻辑。

我认为我需要以某种方式共享或发布流,但我之前使用这些运算符收效甚微,而且不知道如何在不添加第二个订阅的情况下做到这一点。

如何在不订阅两次的情况下共享一个流并以两种不同的方式使用它?我可以创建某种仅在基于订阅的可观察对象时订阅的被动流吗?

【问题讨论】:

    标签: javascript typescript rxjs


    【解决方案1】:

    我最终解决了这个问题(在 @MonkeyMagiic 的帮助下)。

    我必须共享流,以便我可以对数据执行 2 项不同的操作,分别对每个值进行缓冲和处理。问题是这两个流都必须订阅,但我不想订阅服务 - 这应该在控制器中完成。

    解决方案是再次合并 2 个流并忽略缓冲区中的值:

    var intervalStream = Rx.Observable.interval(250)
        .take(8)
        .do( function(value){console.log( "source: " + value );} )
        .shareReplay(1);
    
    var bufferStream = intervalStream.bufferWithCount(3)
        .do( function(values){
          console.log( "do something with buffered values: " + values );
        } )
        .flatMap( function(values){ return Rx.Observable.empty(); } );
    
    var mergeStream = intervalStream.merge( bufferStream );
    
    mergeStream.subscribe(
      function( value ){ console.log( "value received by controller: " + value ); },
      function( error ){ console.log( "error: " + error ); },
      function(){ console.log( "onComplete" ); }
    );
    

    输出:

    "source: 0"
    "value received by controller: 0"
    "source: 1"
    "value received by controller: 1"
    "source: 2"
    "value received by controller: 2"
    "do something with buffered values: 0,1,2"
    "source: 3"
    "value received by controller: 3"
    "source: 4"
    "value received by controller: 4"
    "source: 5"
    "value received by controller: 5"
    "do something with buffered values: 3,4,5"
    "source: 6"
    "value received by controller: 6"
    "source: 7"
    "value received by controller: 7"
    "do something with buffered values: 6,7"
    "onComplete"
    

    JSBin

    【讨论】:

      【解决方案2】:

      我相信您正在寻找的是sharereplay 的组合,幸好RX 有shareReplay(bufferSize)

      var threadStream: Rx.Observable<ICommentThread> = this.youTubeService
                      .getCommentThreadsForChannel()
                      .shareReplay(1); 
      
      getCommentThreadsForChannel(): Rx.Observable<ICommentThread> {
              threadStream
                  .bufferWithCount(100)
                  .scan( ( allItems, currentItem ) => {
                      currentItem.forEach(thread => {
                          allItems.push(thread);
                      });
      
                      console.log( `Save items to local storage: ${allItems.length}` )
      
                      return allItems;
                  }, []  );
      
              return threadStream;
          }
      

      使用多播运营商共享将确保第一个之后的每个附加订阅都不会发出不必要的请求,并且重播是为了确保所有未来的订阅都收到最后一个通知。

      【讨论】:

      • 值得注意的是,shareReplay 目前已从最新的 RxJS 5 测试版中删除。 More info here.
      • 感谢您的回复,但恐怕它不起作用。我认为问题在于从该函数返回时正在订阅线程流,但未订阅 bufferWithCount 和扫描流。正如我之前提到的,我不想在这里订阅,我只希望它在订阅线程流时工作。
      • 我有点迷路了(抱歉),整个 observables 链将作为一个单元订阅,即 {observable1}.{observable2}.{observable3} 将从链向下发送顶部通知 --> observable1 到 observable2。
      • 作为补充,我可以保证您订阅的是 bufferWithCount 和扫描,如果您订阅的是 threadStream。
      • 使用 bufferWithCount(100) 表示您正在等待从 getCommentThreadsForChannel() 释放 100 个通知,这会是您想要的行为吗?你能粘贴那个方法的样子吗?
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-09-10
      • 2016-01-22
      • 2017-11-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-05-11
      相关资源
      最近更新 更多