rxjs 5 publishReplay refCount

Ole*_*llo 17 rxjs rxjs5

我无法弄清楚它是如何publishReplay().refCount()运作的.

例如(https://jsfiddle.net/7o3a45L1/):

var source = Rx.Observable.create(observer =>  {
  console.log("call"); 
  // expensive http request
  observer.next(5);
}).publishReplay().refCount();

subscription1 = source.subscribe({next: (v) => console.log('observerA: ' + v)});
subscription1.unsubscribe();
console.log(""); 

subscription2 = source.subscribe({next: (v) => console.log('observerB: ' + v)});
subscription2.unsubscribe();
console.log(""); 

subscription3 = source.subscribe({next: (v) => console.log('observerC: ' + v)});
subscription3.unsubscribe();
console.log(""); 

subscription4 = source.subscribe({next: (v) => console.log('observerD: ' + v)});
subscription4.unsubscribe();
Run Code Online (Sandbox Code Playgroud)

给出以下结果:

致电观察员:5

observerB:5调用observerB:5

observerC:5 observerC:5 call observerC:5

观察者:5观察者:5观察者:5,观察者观察者:5

1)为什么observerB,C和D被多次调用?

2)为什么"call"打印在每一行而不是行的开头?

此外,如果我打电话publishReplay(1).refCount(),它会分别调用observerB,C和D 2次.

我期望的是每个新观察者只接收一次值5,而"呼叫"只打印一次.

Mar*_*ten 24

publishReplay(x).refCount() 合并执行以下操作:

  • 它创造了一个ReplaySubject重放x排放的重放.如果未定义x,则重放完整的流.
  • 它ReplaySubject使用refCount()运算符使此多播兼容.这导致并发订阅接收相同的排放.

您的示例包含一些问题,使它们如何一起工作.请参阅以下修订的代码段:

var state = 5
var realSource = Rx.Observable.create(observer =>  {
  console.log("creating expensive HTTP-based emission"); 
  observer.next(state++);
//  observer.complete();
  
  return () => {
    console.log('unsubscribing from source')
  }
});


var source = Rx.Observable.of('')
  .do(() => console.log('stream subscribed'))
  .ignoreElements()
  .concat(realSource)
.do(null, null, () => console.log('stream completed'))
.publishReplay()
.refCount()
;
    
subscription1 = source.subscribe({next: (v) => console.log('observerA: ' + v)});
subscription1.unsubscribe();
 
subscription2 = source.subscribe(v => console.log('observerB: ' + v));
subscription2.unsubscribe();
    
subscription3 = source.subscribe(v => console.log('observerC: ' + v));
subscription3.unsubscribe();
    
subscription4 = source.subscribe(v => console.log('observerD: ' + v));
 
Run Code Online (Sandbox Code Playgroud)
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.1.0/Rx.js"></script>
Run Code Online (Sandbox Code Playgroud)

在运行此代码段时,我们可以清楚地看到它没有为Observer D发出重复值,实际上它为每个订阅创建了新的排放量.怎么会?

在下一次订阅发生之前,每个订阅都会被取消订阅.这有效地使refCount减少到零,没有进行多播.

问题在于realSource流未完成.因为我们不是多播,所以下一个订户realSource通过ReplaySubject 获得一个新的实例,并且新的排放量先于已经排放的排放量.

因此,为了修复您的流,多次调用昂贵的HTTP请求,您必须完成流,以便publishReplay知道它不需要重新订阅.


ols*_*lsn 7

通常:refCount只要有至少1个订户,流即热/共享的方式 - 但是,当没有订户时,它正在重置/冷.

这意味着如果您想要绝对确定不会执行任何操作,则不应使用refCount()简单connect的流来设置它.

作为补充说明:如果您在observer.complete()之后添加,observer.next(5);您还将获得预期的结果.


旁注:你真的需要在Obervable这里创建自己的自定义吗?在95%的情况下,现有的运算符足以满足给定的用例.


mar*_*tin 5

发生这种情况是因为您正在使用publishReplay(). 它在内部创建一个实例,ReplaySubject用于存储所有通过的值。

由于您使用的Observable.create是发出单个值的地方,因此每次调用时source.subscribe(...)都会将一个值附加到ReplaySubject.

您不会call在每一行的开头打印,因为它是ReplaySubject在您订阅时首先发出缓冲区的人,然后再订阅其源:

实现细节见:

使用publishReplay(1). 首先它发出缓冲的项ReplaySubject,然后再发出另一个项observer.next(5);