RxJS的队列运算符

Arl*_*ler 5 rxjs reactivex

在RxJS中是否有一个操作符允许我缓冲项目并在信号可观察到的时候逐个发出它们?类似于bufferWhen,但是不是在每个信号上转储整个缓冲区,而是每个信号转储一定数量的缓冲区.它甚至可以转储信号可观察到的数字.

Input observable:  >--a--b--c--d--|
Signal observable: >------1---1-1-|
Count in buffer:   !--1--21-2-121-|
Output observable: >------a---b-c-|
Run Code Online (Sandbox Code Playgroud)

car*_*ant 6

是的,你可以zip用来做你想做的事:

const input = Rx.Observable.from(["a", "b", "c", "d", "e"]);
const signal = new Rx.Subject();
const output = Rx.Observable.zip(input, signal, (i, s) => i);
output.subscribe(value => console.log(value));
signal.next(1);
signal.next(1);
signal.next(1);
Run Code Online (Sandbox Code Playgroud)
.as-console-wrapper { max-height: 100% !important; top: 0; }
Run Code Online (Sandbox Code Playgroud)
<script src="https://unpkg.com/rxjs@5/bundles/Rx.min.js"></script>
Run Code Online (Sandbox Code Playgroud)

实际上,zip这个与缓冲有关的GitHub问题中用作示例.

如果要使用信号的发射值来确定要释放多少缓冲值,可以执行以下操作:

const input = Rx.Observable.from(["a", "b", "c", "d", "e"]);
const signal = new Rx.Subject();
const output = Rx.Observable.zip(
  input,
  signal.concatMap(count => Rx.Observable.range(0, count)),
  (i, s) => i
);
output.subscribe(value => console.log(value));
signal.next(1);
signal.next(2);
signal.next(1);
Run Code Online (Sandbox Code Playgroud)
.as-console-wrapper { max-height: 100% !important; top: 0; }
Run Code Online (Sandbox Code Playgroud)
<script src="https://unpkg.com/rxjs@5/bundles/Rx.min.js"></script>
Run Code Online (Sandbox Code Playgroud)