Replay
обеспечить, чтобы все наблюдатели увидели одну и ту же последовательность выпущенных элементов, даже если они подпишутся после того, как Observable начал выпускать элементы
A подключаемый Observable похож на обычный Observable, за исключением того, что он не начинает выпускать элементы при подписке, а только тогда, когда к нему применяется оператор Connect. Таким образом, вы можете побудить Observable начать выпускать элементы в нужное время.
Если вы примените оператор Replay к Observable перед преобразованием его в подключаемый Observable, полученный подключаемый Observable всегда будет выпускать ту же полную последовательность всем будущим наблюдателям, даже тем наблюдателям, которые подписываются после того, как подключаемый Observable начал выпускать элементы другим подписанным наблюдателям.
См. также
Информация для определённого языка
RxGroovy replay cache
В RxGroovy существует множество оператора replay, который возвращает соединяемый Observable. Вам необходимо опубликовать этот соединяемый Observable до того, как подписчики смогут к нему подключиться, а затем подключиться к нему, чтобы наблюдать за его выбросами.
Разновидности этого множества оператора replay позволяют задать максимальный размер буфера, чтобы ограничить количество элементов, которые replay будут буферизовать и повторно отправлять последующим подписчикам, и/или установить подвижное временное окно, которое определяет, когда отправленные элементы становятся слишком старыми, чтобы их буферизовать и повторно отправлять.
- Javadoc:
replay() - Javadoc:
replay(Scheduler) - Javadoc:
replay(int) - Javadoc:
replay(int,Scheduler) - Javadoc:
replay(long,TimeUnit) - Javadoc:
replay(long,TimeUnit,Scheduler) - Javadoc:
replay(int,long,TimeUnit) - Javadoc:
replay(int,long,TimeUnit,Scheduler)
Также существует множество оператора replay, который возвращает обычный Observable. Эти варианты принимают в качестве параметра функцию преобразования; эта функция принимает в качестве параметра элемент, испускаемый исходным Observable, и возвращает элемент, который должен быть испущен результирующим Observable. Так что на самом деле, этот оператор не повторяет исходный Observable, а вместо этого повторяет исходный Observable, *преобразованный* этой функцией.
Разновидности этого множества оператора replay позволяют задать максимальный размер буфера, чтобы ограничить количество элементов, которые replay будут буферизовать и повторно отправлять последующим подписчикам, и/или установить подвижное временное окно, которое определяет, когда отправленные элементы становятся слишком старыми, чтобы их буферизовать и повторно отправлять.
- Javadoc:
replay(Func1) - Javadoc:
replay(Func1,Scheduler) - Javadoc:
replay(Func1,int) - Javadoc:
replay(Func1,int,Scheduler) - Javadoc:
replay(Func1,long,TimeUnit) - Javadoc:
replay(Func1,long,TimeUnit,Scheduler) - Javadoc:
replay(Func1,int,long,TimeUnit) - Javadoc:
replay(Func1,int,long,TimeUnit,Scheduler)
RxJava 1․x cache replay
В RxJava существует множество оператора replay, который возвращает соединяемый Observable. Вам необходимо опубликовать этот соединяемый Observable до того, как подписчики смогут к нему подключиться, а затем подключиться к нему, чтобы наблюдать за его выбросами.
Разновидности этого множества оператора replay позволяют задать максимальный размер буфера, чтобы ограничить количество элементов, которые replay будут буферизовать и повторно отправлять последующим подписчикам, и/или установить подвижное временное окно, которое определяет, когда отправленные элементы становятся слишком старыми, чтобы их буферизовать и повторно отправлять.
- Javadoc:
replay() - Javadoc:
replay(Scheduler) - Javadoc:
replay(int) - Javadoc:
replay(int,Scheduler) - Javadoc:
replay(long,TimeUnit) - Javadoc:
replay(long,TimeUnit,Scheduler) - Javadoc:
replay(int,long,TimeUnit) - Javadoc:
replay(int,long,TimeUnit,Scheduler)
Также существует множество оператора replay, который возвращает обычный Observable. Эти варианты принимают в качестве параметра функцию преобразования; эта функция принимает в качестве параметра элемент, испускаемый исходным Observable, и возвращает элемент, который должен быть испущен результирующим Observable. Так что на самом деле, этот оператор не повторяет исходный Observable, а вместо этого повторяет исходный Observable, *преобразованный* этой функцией.
Разновидности этого множества оператора replay позволяют задать максимальный размер буфера, чтобы ограничить количество элементов, которые replay будут буферизовать и повторно отправлять последующим подписчикам, и/или установить подвижное временное окно, которое определяет, когда отправленные элементы становятся слишком старыми, чтобы их буферизовать и повторно отправлять.
- Javadoc:
replay(Func1) - Javadoc:
replay(Func1,Scheduler) - Javadoc:
replay(Func1,int) - Javadoc:
replay(Func1,int,Scheduler) - Javadoc:
replay(Func1,long,TimeUnit) - Javadoc:
replay(Func1,long,TimeUnit,Scheduler) - Javadoc:
replay(Func1,int,long,TimeUnit) - Javadoc:
replay(Func1,int,long,TimeUnit,Scheduler)
RxJS replay shareReplay
В RxJs оператор replay принимает четыре необязательных параметра и возвращает обычный Observable:
selector- функция преобразования, которая принимает элемент, испускаемый исходным Observable, в качестве параметра и возвращает элемент, который должен быть испущен результирующим Observable
bufferSize- максимальное количество элементов для буферизации и повторной отправки последующим подписчикам
window- возраст в миллисекундах, по истечении которого элементы в этом буфере могут быть удалены без отправки последующим подписчикам
scheduler- Scheduler, на котором будет работать этот оператор
Пример кода
var interval = Rx.Observable.interval(1000);
var source = interval
.take(2)
.do(function (x) {
console.log('Side effect');
});
var published = source
.replay(function (x) {
return x.take(2).repeat(2);
}, 3);
published.subscribe(createObserver('SourceA'));
published.subscribe(createObserver('SourceB'));
function createObserver(tag) {
return Rx.Observer.create(
function (x) { console.log('Next: ' + tag + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); });
}
Side effect Next: SourceA0 Side effect Next: SourceB0 Side effect Next: SourceA1 Next: SourceA0 Next: SourceA1 Completed Side effect Next: SourceB1 Next: SourceB0 Next: SourceB1 Completed
Также существует оператор shareReplay, который отслеживает количество подписчиков и отключается от исходного Observable, когда это число падает до нуля. Оператор shareReplay принимает три необязательных параметра и возвращает обычный Observable:
bufferSize- максимальное количество элементов для буферизации и повторной отправки последующим подписчикам
window- возраст в миллисекундах, по истечении которого элементы в этом буфере могут быть удалены без отправки последующим подписчикам
scheduler- Scheduler, на котором будет работать этот оператор
Пример кода
var interval = Rx.Observable.interval(1000);
var source = interval
.take(4)
.doAction(function (x) {
console.log('Side effect');
});
var published = source
.shareReplay(3);
published.subscribe(createObserver('SourceA'));
published.subscribe(createObserver('SourceB'));
// Creating a third subscription after the previous two subscriptions have
// completed. Notice that no side effects result from this subscription,
// because the notifications are cached and replayed.
Rx.Observable
.return(true)
.delay(6000)
.flatMap(published)
.subscribe(createObserver('SourceC'));
function createObserver(tag) {
return Rx.Observer.create(
function (x) { console.log('Next: ' + tag + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); });
}
Side effect Next: SourceA0 Next: SourceB0 Side effect Next: SourceA1 Next: SourceB1 Side effect Next: SourceA2 Next: SourceB2 Side effect Next: SourceA3 Next: SourceB3 Completed Completed Next: SourceC1 Next: SourceC2 Next: SourceC3 Completed
Операторы replay и shareReplay находятся в следующих дистрибутивах:
rx.all.jsrx.all.compat.js-
rx.binding.js(требуетrx.jsилиrx.compat.js) rx.lite.jsrx.lite.compat.js
RxPHP replay shareReplay
RxPHP реализует этот оператор как replay.
Возвращает последовательность наблюдаемых значений, которая является результатом вызова селектора над соединяемой последовательностью наблюдаемых значений, которая разделяет одну подписку на исходную последовательность, повторно воспроизводящую уведомления, ограниченные максимальной длительностью для буфера воспроизведения. Этот оператор является специализацией Multicast, использующей ReplaySubject.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/replay/replay.php
$interval = \Rx\Observable::interval(1000);
$source = $interval
->take(2)
->doOnNext(function ($x) {
echo $x, ' something', PHP_EOL;
echo 'Side effect', PHP_EOL;
});
$published = $source
->replay(function (\Rx\Observable $x) {
return $x->take(2)->repeat(2);
}, 3);
$published->subscribe($createStdoutObserver('SourceA '));
$published->subscribe($createStdoutObserver('SourceB ')); 0 something Side effect 0 something Side effect SourceA Next value: 0 SourceB Next value: 0 SourceA Next value: 0 SourceB Next value: 0 SourceA Next value: 0 SourceB Next value: 0 SourceA Next value: 0 SourceA Complete! SourceB Next value: 0 SourceB Complete! 1 something Side effect 1 something Side effect
RxPHP также имеет оператор shareReplay.
Возвращает последовательность наблюдаемых значений, которая разделяет одну подписку на исходную последовательность, повторно воспроизводящую уведомления, ограниченные максимальной длительностью для буфера воспроизведения. Этот оператор является специализацией replay, который создаёт подписку, когда количество наблюдателей меняется с нуля на единицу, затем разделяет эту подписку со всеми последующими наблюдателями, пока количество наблюдателей не вернётся к нулю, в этот момент подписка удаляется.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/share/shareReplay.php
$interval = Rx\Observable::interval(1000);
$source = $interval
->take(4)
->doOnNext(function ($x) {
echo 'Side effect', PHP_EOL;
});
$published = $source
->shareReplay(3);
$published->subscribe($createStdoutObserver('SourceA '));
$published->subscribe($createStdoutObserver('SourceB '));
Rx\Observable
::of(true)
->concatMapTo(\Rx\Observable::timer(6000))
->flatMap(function () use ($published) {
return $published;
})
->subscribe($createStdoutObserver('SourceC ')); Side effect SourceA Next value: 0 SourceB Next value: 0 Side effect SourceA Next value: 1 SourceB Next value: 1 Side effect SourceA Next value: 2 SourceB Next value: 2 Side effect SourceA Next value: 3 SourceB Next value: 3 SourceA Complete! SourceB Complete! SourceC Next value: 1 SourceC Next value: 2 SourceC Next value: 3 SourceC Complete!
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/replay.html