Spec-Zone.ru › ReactiveX

Образец

выдавать самые последние элементы, испущенные Observable, в периодические интервалы времени
Открыть интерактивную диаграмму на rxmarbles.com

Оператор Sample периодически просматривает Observable и выдаёт тот элемент, который был последним выпущенным с момента предыдущего опроса.

В некоторых реализациях также есть оператор ThrottleFirst, который похож, но выдаёт не последний элемент в периоде опроса, а первый элемент, который был выпущен в этот период.

См. также

  • Операторы, связанные с обратной пропускной способностью
  • Debounce
  • Window
  • Введение в Rx: Sample
  • RxMarbles: sample
  • 101 Примеры Rx: Sample — Простой

Информация, специфичная для языка

RxGroovy sample throttleFirst throttleLast

RxGroovy реализует этот оператор как sample и throttleLast.

Обратите внимание, что если исходный Observable не выпустил никаких элементов с последнего момента опроса, Observable, полученный в результате этого оператора, не выпустит ни одного элемента в этот период опроса.

sample

Один вариант 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 по вашему выбору в качестве третьего параметра.

  • Javadoc: sample(long,TimeUnit) и throttleLast(long,TimeUnit)
  • Javadoc: sample(long,TimeUnit,Scheduler) и throttleLast(long,TimeUnit,Scheduler)
sample

Также есть вариант sample (у которого нет псевдонима throttleLast) опроса исходного Observable каждый раз, когда второй Observable выдаёт элемент (или когда он завершается). Вы передаёте этот второй Observable в качестве параметра sample.

Этот вариант sample по умолчанию не работает с каким-либо конкретным Scheduler.

  • Javadoc: sample(Observable)
throttleFirst

Также есть оператор 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 по вашему выбору в качестве третьего параметра.

  • throttleFirst(long,TimeUnit)
  • throttleFirst(long,TimeUnit,Scheduler)

RxJava 1․x sample throttleFirst throttleLast

RxJava реализует этот оператор как sample и throttleLast.

Обратите внимание, что если исходный Observable не выпустил никаких элементов с последнего момента опроса, Observable, полученный в результате этого оператора, не выпустит ни одного элемента в этот период опроса.

sample

Один вариант sample (или его псевдоним, throttleLast) делает опрос с периодическим интервалом времени, который вы выбираете, передавая TimeUnit и количество таких единиц в качестве параметров sample.

Этот вариант sample по умолчанию работает с computation Scheduler, но вы можете необязательно указать Scheduler по вашему выбору в качестве третьего параметра.

  • Javadoc: sample(long,TimeUnit) и throttleLast(long,TimeUnit)
  • Javadoc: sample(long,TimeUnit,Scheduler) и throttleLast(long,TimeUnit,Scheduler)
sample

Также есть вариант sample (у которого нет псевдонима throttleLast) опроса исходного Observable каждый раз, когда второй Observable выдаёт элемент (или когда он завершается). Вы передаёте этот второй Observable в качестве параметра sample.

Этот вариант sample по умолчанию не работает с каким-либо конкретным Scheduler.

  • Javadoc: sample(Observable)
throttleFirst

Также есть оператор throttleFirst, который отличается от throttleLast/sample тем, что выдаёт первый элемент, выпущенный исходным Observable в каждом периоде опроса, а не последний выпущенный элемент.

throttleFirst по умолчанию работает с computation Scheduler, но вы можете необязательно указать Scheduler по вашему выбору в качестве третьего параметра.

  • throttleFirst(long,TimeUnit)
  • throttleFirst(long,TimeUnit,Scheduler)

RxJS sample throttleFirst

RxJS реализует этот оператор с двумя вариантами sample.

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
sample

Второй вариант принимает в качестве параметра 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

Также есть оператор 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.js
  • rx.all.compat.js
  • rx.time.js (требует rx.js или rx.compat.js)
  • rx.lite.js
  • rx.lite.compat.js

© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/sample.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API