javascript – RxJs:zip运算符的有损形式

考虑使用 zip运算符将两个无限的Observable压缩在一起,其中一个Observable发出的项目频率是另一个的两倍.
当前的实现是无损耗的,即如果我让这些Observables发射一小时然后我在它们的发射速率之间切换,第一个Observable将最终赶上另一个.
随着缓冲区越来越大,这将导致内存爆炸.
如果第一个observable将发出几个小时的项目而第二个将在最后发出一个项目,则会发生同样的情况.

如何为此操作符实现有损行为?我只是想随时从两个流中获得排放,而我不关心我错过的更快流量的排放量.

澄清:

>我试图在这里解决的主要问题是由于zip运算符的无损性质导致的内存爆炸.
>我希望随时都能从两个流中获得发射,即使两个流每次都发出相同的值

例:

Stream1: 1 2    3 4    5 6 7                
Stream2:     10     20       30 40 50 60 70

常规zip将产生以下输出

[1,10]
[2,20]
[3,30]
[4,40]
[5,50]
[6,60]
[7,70]
const Observable = Rx.Observable;
const Subject = Rx.Subject;


const s1 = new Subject();
const s2 = new Subject();

Observable.zip(s1,s2).subscribe(console.log);

s1.next(1); s1.next(2); s2.next(10); s1.next(3); s1.next(4); s2.next(20); s1.next(5); s1.next(6); s1.next(7); s2.next(30); 
 
s2.next(40); s2.next(50); s2.next(60); s2.next(70);
<script src="https://unpkg.com/@reactivex/rxjs@5.0.3/dist/global/Rx.js"></script>

我希望它产生的输出

[1,10]
[3,20]
[5,30]

说明:
有损zip运算符是zip,缓冲区大小为1.这意味着它只保留首先发出的流中的第一个项目,并且将丢弃所有其余项目(在第一个项目和第二个流的第一个发射之间到达的项目).因此,示例中发生的情况如下:stream1发出1,有损zip会“记住”它并忽略stream1上的所有项目,直到stream2发出.第一次发射的stream2是10,所以stream1松散2.在相互发射(第一次发射有损拉链)后,它重新开始:“记住”3,“松散”4,发出[3,20].然后重新开始:“记住”5,“松散”6和7,发出[5,30].然后重新开始:“记住”40,“松散”50,60,70并等待stream1上的下一个项目.

例2:

Stream1: 1 2 3 ... 100000000000
Stream2:                        a

在这种情况下,常规zip操作符会爆炸内存.
我不想要它.

摘要
基本上我希望有损zip运算符只记住前一次相互发射后流1发出的第一个值,并在流2赶上流1时发出.然后重复.

解决方法

以下内容将为您提供所需的行为:
Observable.zip(s1.take(1),s2.take(1)).repeat()

在RxJs 5.5管道语法中:

zip(s1.pipe(take(1)),s2.pipe(take(1))).pipe(repeat());
const s1 = new Rx.Subject();
const s2 = new Rx.Subject();

Rx.Observable.zip(s1.take(1),s2.take(1)).repeat()
    .subscribe(console.log);

s1.next(1); s1.next(2); s2.next(10); s1.next(3); s1.next(4); s2.next(20); s1.next(5); s1.next(6); s1.next(7); s2.next(30);  
s2.next(40); s2.next(50); s2.next(60); s2.next(70);
<script src="https://unpkg.com/@reactivex/rxjs@5.0.3/dist/global/Rx.js"></script>

说明:

>重复运算符(在其当前实现中)在后者完成时重新订阅可观察到的源,即在该特定情况下,它在每次相互发射时重新订阅以压缩.
> zip结合了两个observable并等待它们两个发出. combineLatest也会这样做,因为take(1)并不重要
> take(1)实际上处理内存爆炸并定义有损行为

如果你想在相互发射时从每个流中获取最后一个而不是第一个值,请使用:

Observable.combineLatest(s1,s2).take(1).repeat()

在RxJs 5.5管道语法中:

combineLatest(s1.pipe(take(1)),s2.pipe(take(1))).pipe(repeat());
const s1 = new Rx.Subject();
const s2 = new Rx.Subject();

Rx.Observable.combineLatest(s1,s2).take(1).repeat()
    .subscribe(console.log);

s1.next(1); s1.next(2); s2.next(10); s1.next(3); s1.next(4); s2.next(20); s1.next(5); s1.next(6); s1.next(7); s2.next(30);  
s2.next(40); s2.next(50); s2.next(60); s2.next(70);
<script src="https://unpkg.com/@reactivex/rxjs@5.0.3/dist/global/Rx.js"></script>

相关文章

事件冒泡和事件捕获 起因:今天在封装一个bind函数的时候,发现el.addEventListener函数支持第三个参数...
js小数运算会出现精度问题 js number类型 JS 数字类型只有number类型,number类型相当于其他强类型语言...
什么是跨域 跨域 : 广义的跨域包含一下内容 : 1.资源跳转(链接跳转,重定向跳转,表单提交) 2.资源...
@ &quot;TOC&quot; 常见对base64的认知(不完全正确) 首先对base64常见的认知,也是须知的必须有...
搞懂:MVVM模式和Vue中的MVVM模式 MVVM MVVM : 的缩写,说都能直接说出来 :模型, :视图, :视图模...
首先我们需要一个html代码的框架如下: 我们的目的是实现ul中的内容进行横向的一点一点滚动。ul中的内容...