Образец
выдавать самые последние элементы, испущенные Observable, в периодические интервалы времени
Оператор Sample периодически просматривает Observable и выдаёт тот элемент, который был последним выпущенным с момента предыдущего опроса.
В некоторых реализациях также есть оператор ThrottleFirst, который похож, но выдаёт не последний элемент в периоде опроса, а первый элемент, который был выпущен в этот период.
См. также
- Операторы, связанные с обратной пропускной способностью
- Debounce
- Window
- Введение в Rx: Sample
- RxMarbles:
sample - 101 Примеры Rx: Sample — Простой
Информация, специфичная для языка
RxGroovy sample throttleFirst throttleLast
RxGroovy реализует этот оператор как sample и throttleLast.
Обратите внимание, что если исходный Observable не выпустил никаких элементов с последнего момента опроса, Observable, полученный в результате этого оператора, не выпустит ни одного элемента в этот период опроса.
Один вариант sample (или его псевдоним, throttleLast) делает опрос с периодическим интервалом времени, который вы выбираете, передавая TimeUnit и количество таких единиц в качестве параметров sample.
Следующий код создаёт Observable, который выдаёт числа от одного до миллиона, а затем опрошивает этот Observable каждые десять миллисекунд, чтобы увидеть, какое число он выдаёт в этот момент.
Пример кода
def numbers = Observable.range( 1, 1000000 );
numbers.sample(10, java.util.concurrent.TimeUnit.MILLISECONDS).subscribe(
{ println(it); }, // onNext
{ println("Error: " + it.getMessage()); }, // onError
{ println("Sequence complete"); } // onCompleted
); 339707 547810 891282 Sequence complete
Этот вариант sample по умолчанию работает с computation Scheduler, но вы можете необязательно указать Scheduler по вашему выбору в качестве третьего параметра.
Также есть вариант sample (у которого нет псевдонима throttleLast) опроса исходного Observable каждый раз, когда второй Observable выдаёт элемент (или когда он завершается). Вы передаёте этот второй Observable в качестве параметра sample.
Этот вариант sample по умолчанию не работает с каким-либо конкретным Scheduler.
- Javadoc:
sample(Observable)
Также есть оператор throttleFirst, который отличается от throttleLast/sample тем, что выдаёт первый элемент, выпущенный исходным Observable в каждом периоде опроса, а не последний выпущенный элемент.
Пример кода
Scheduler s = new TestScheduler();
PublishSubject<Integer> o = PublishSubject.create();
o.throttleFirst(500, TimeUnit.MILLISECONDS, s).subscribe(
{ println(it); }, // onNext
{ println("Error: " + it.getMessage()); }, // onError
{ println("Sequence complete"); } // onCompleted
);
// send events with simulated time increments
s.advanceTimeTo(0, TimeUnit.MILLISECONDS);
o.onNext(1); // deliver
o.onNext(2); // skip
s.advanceTimeTo(501, TimeUnit.MILLISECONDS);
o.onNext(3); // deliver
s.advanceTimeTo(600, TimeUnit.MILLISECONDS);
o.onNext(4); // skip
s.advanceTimeTo(700, TimeUnit.MILLISECONDS);
o.onNext(5); // skip
o.onNext(6); // skip
s.advanceTimeTo(1001, TimeUnit.MILLISECONDS);
o.onNext(7); // deliver
s.advanceTimeTo(1501, TimeUnit.MILLISECONDS);
o.onCompleted(); 1 3 7 Sequence complete
throttleFirst по умолчанию работает с computation Scheduler, но вы можете необязательно указать Scheduler по вашему выбору в качестве третьего параметра.
RxJava 1․x sample throttleFirst throttleLast
RxJava реализует этот оператор как sample и throttleLast.
Обратите внимание, что если исходный Observable не выпустил никаких элементов с последнего момента опроса, Observable, полученный в результате этого оператора, не выпустит ни одного элемента в этот период опроса.
Один вариант sample (или его псевдоним, throttleLast) делает опрос с периодическим интервалом времени, который вы выбираете, передавая TimeUnit и количество таких единиц в качестве параметров sample.
Этот вариант sample по умолчанию работает с computation Scheduler, но вы можете необязательно указать Scheduler по вашему выбору в качестве третьего параметра.
Также есть вариант sample (у которого нет псевдонима throttleLast) опроса исходного Observable каждый раз, когда второй Observable выдаёт элемент (или когда он завершается). Вы передаёте этот второй Observable в качестве параметра sample.
Этот вариант sample по умолчанию не работает с каким-либо конкретным Scheduler.
- Javadoc:
sample(Observable)
Также есть оператор throttleFirst, который отличается от throttleLast/sample тем, что выдаёт первый элемент, выпущенный исходным Observable в каждом периоде опроса, а не последний выпущенный элемент.
throttleFirst по умолчанию работает с computation Scheduler, но вы можете необязательно указать Scheduler по вашему выбору в качестве третьего параметра.
RxJS sample throttleFirst
RxJS реализует этот оператор с двумя вариантами sample.
Первый вариант принимает в качестве параметра периодичность, определённую как целое число миллисекунд, и периодически опрашивает исходный Observable с этой частотой.
Пример кода
var source = Rx.Observable.interval(1000)
.sample(5000)
.take(2);
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }); Next: 3 Next: 8 Completed
Второй вариант принимает в качестве параметра Observable, и опрашивает исходный Observable всякий раз, когда этот второй Observable выдаёт элемент.
Пример кода
var source = Rx.Observable.interval(1000)
.sample(Rx.Observable.interval(5000))
.take(2);
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }); Next: 3 Next: 8 Completed
Также есть оператор throttleFirst, который отличается от sample тем, что выдаёт первый элемент, выпущенный исходным Observable в каждом периоде опроса, а не последний выпущенный элемент.
У него нет варианта, использующего выдачи второго Observable для регулирования периодичности опроса.
Пример кода
var times = [
{ value: 0, time: 100 },
{ value: 1, time: 600 },
{ value: 2, time: 400 },
{ value: 3, time: 900 },
{ value: 4, time: 200 }
];
// Delay each item by time and project value;
var source = Rx.Observable.from(times)
.flatMap(function (item) {
return Rx.Observable
.of(item.value)
.delay(item.time);
})
.throttleFirst(300 /* ms */);
var subscription = source.subscribe(
function (x) { console.log('Next: %s', x); },
function (err) { console.log('Error: %s', err); },
function () { console.log('Completed'); }); Next: 0 Next: 2 Next: 3 Completed
sample и throttleFirst по умолчанию работают с timeout Scheduler. Они находятся в следующих распределениях:
rx.all.jsrx.all.compat.js-
rx.time.js(требуетrx.jsилиrx.compat.js) rx.lite.jsrx.lite.compat.js
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/sample.html