RXJS Observable的冷,热和Subject

一、Observable的冷和热

Observable 热:直播。所有的观察者,无论进来的早还是晚,看到的是同样内容的同样进度,订阅的时候得到的都是最新时刻发送的值。

Observable 冷:点播。 新的订阅者每次从头开始。

冷的Observable例子:

一开始有个订阅者,

两秒后又有个订阅者,这两个序列按照自己的节奏走的,不同步。每个流进行都会从interval的0开始。

1console.log('RxJS included?', !!Rx); 2 3const count$ = Rx.Observable.interval(1000).take(5); 4const sub1 = count$.subscribe((val)=>{ 5 console.log(val); 6}); 7 8setTimeout(function(){ 9 const sub2 = count$.subscribe((val)=>{ 10 console.log(val); 11}); 12},2000);

热的Observable例子

第二个订阅者直接从2开始起,跟第一个订阅者看到的内容是一样的。

1const count$ = Rx.Observable.interval(1000).take(5).share(); 2const sub1 = count$.subscribe((val)=>{ 3 console.log(val); 4}); 5 6setTimeout(function(){ 7 const sub2 = count$.subscribe((val)=>{ 8 console.log(val); 9}); 10},2000);

二、Subject

Subject即是观察者Observer,也是被观察对象Observable,同时实现了这两个接口。

意味着

  • 一方面它可以作为流的组成的一方,输出的一方。
  • 另一方面,它可以作为流的观察一方,接收一方。

Subject分为ReplaySubject和BehaviorSubject。

ReplaySubject:这种Subject会保留最新的n个值

BehaviorSubject:是ReplaySubject的特殊形式。 保留最新的一个值

1、subscribe的等价写法

subscribe 后面写的一个函数,相当于语法糖,快捷方式,临时创建冷一个observer对象。

默认情况应该是传入一个observer对象

1console.log('RxJS included?', !!Rx); 2 3 4const counter$ = Rx.Observable.interval(1000).take(5); 5 6const subject = new Rx.Subject(); 7 8const observer1 = { 9 next: (val)=>{console.log('1: ' +val);}, 10 error: (err)=>{console.log('ERROR>> 1:'+ err);}, 11 complete: ()=>{console.log('1 is complete');} 12} 13 14 15const observer2 = { 16 next: (val)=>{console.log('2: ' +val);}, 17 error: (err)=>{console.log('ERROR>> 2:'+ err);}, 18 complete: ()=>{console.log('2 is complete');} 19} 20 21//等价写法 22counter$.subscribe(val =>{console.log(val);}); 23counter$.subscribe(observer2);

2、两个observer ,两次subscribe

1console.log('RxJS included?', !!Rx); 2 3 4const counter$ = Rx.Observable.interval(1000).take(5); 5 6const subject = new Rx.Subject(); 7 8const observer1 = { 9 next: (val)=>{console.log('1: ' +val);}, 10 error: (err)=>{console.log('ERROR>> 1:'+ err);}, 11 complete: ()=>{console.log('1 is complete');} 12} 13 14 15const observer2 = { 16 next: (val)=>{console.log('2: ' +val);}, 17 error: (err)=>{console.log('ERROR>> 2:'+ err);}, 18 complete: ()=>{console.log('2 is complete');} 19} 20 21counter$.subscribe(observer1); 22 23setTimeout(function(){ 24 counter$.subscribe(observer2); 25},2000);

View Code

 

问题:需要在两处执行subscribe,很多情况下是这样的,定义好这些序列应该在什么时候被触发,我执行执行一句subscribe(),两个序列都会这么执行。这种情况下就需要用subject()。

3、subject

subject即使observable,因为它可以subscribe observer。

也是observer,因为它可以被observable subscribe。

1console.log('RxJS included?', !!Rx); 2 3 4const counter$ = Rx.Observable.interval(1000).take(5); 5 6const subject = new Rx.Subject(); 7 8 9const observer1 = { 10 next: (val)=>{console.log('1: ' +val);}, 11 error: (err)=>{console.log('ERROR>> 1:'+ err);}, 12 complete: ()=>{console.log('1 is complete');} 13} 14 15 16const observer2 = { 17 next: (val)=>{console.log('2: ' +val);}, 18 error: (err)=>{console.log('ERROR>> 2:'+ err);}, 19 complete: ()=>{console.log('2 is complete');} 20} 21 22//不再用counter$去subscribe,用subject去subscribe, 23subject.subscribe(observer1); 24 25setTimeout(function(){ 26 subject.subscribe(observer2); 27},2000); 28 29//定义好两边后,用counter$去subscribe 30counter$.subscribe(subject);

View Code

用一句执行counter$.subscribe(subject),把定义好的序列,包括等待2秒的序列全部完成了。

4,subject是一个hot observable

往流里推送新值

 第二个拿不到新值,因为第二个流订阅的时候,两个新值已经过去了。

5,ReplaySubject

replay把过去发生的事件进行重播。

ReplaySubject(2)把过去的2个事件进行重播。这样observer1 subscribe的时候就可以看到10和11。

6、BehaviorSubject只记住最新的值

总有一个最新值,总记住上一次的最新值

1console.log('RxJS included?', !!Rx); 2 3 4const counter$ = Rx.Observable.interval(1000).take(5); 5 6const subject = new Rx.BehaviorSubject(); 7 8 9subject.next(10); 10subject.next(11); 11const observer1 = { 12 next: (val)=>{console.log('1: ' +val);}, 13 error: (err)=>{console.log('ERROR>> 1:'+ err);}, 14 complete: ()=>{console.log('1 is complete');} 15} 16 17 18const observer2 = { 19 next: (val)=>{console.log('2: ' +val);}, 20 error: (err)=>{console.log('ERROR>> 2:'+ err);}, 21 complete: ()=>{console.log('2 is complete');} 22} 23 24 25//不再用counter$去subscribe,用subject去subscribe, 26subject.subscribe(observer1); 27 28setTimeout(function(){ 29 subject.subscribe(observer2); 30},2000); 31 32//定义好两边后,用counter$去subscribe 33counter$.subscribe(subject);

View Code

取值的时候,会取得到最新的data,尽管在取值的时候也就是subscribre的时候值已经发射完了,尽管时机已经错失了还是能够得到它上一次发射之后的最新的一个值。

三、Angular中对Rx的支持

大量内置Observable支持:如Http,ReactiveForms,Route等。

Async Pipe是什么?有什么用?

Observable需要subscribe 一下,成员数组变量等于Observable得到的值。

使用Async Pipe可以直接使用Observable,还不用去取消订阅。

memberResults$: Observable<User[]>; 

本文作者starof,因知识本身在变化,作者也在不断学习成长,文章内容也不定时更新,为避免误导读者,方便追根溯源,请诸位转载注明出处:https://www.cnblogs.com/starof/p/10505617.html 有问题欢迎与我讨论,共同进步。

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

swap空间的增减方法

(1)增大swap空间去激活swap交换区:swapoff v /dev/vg00/lvswap扩展交换lv:lvextend L 10G /dev/vg00/lvswap重新生成swap交换区:mkswap /dev/vg00/lvswap激活新生成的交换区:swapon v /dev/vg00/lvswap