RefCount
создайте Connectable Observable, ведя себя как обычное Observable
Connectable Observable напоминает обычное Observable, за исключением того, что оно не начинает испускать элементы при подписке, а только когда к нему применяется оператор Connect. Таким образом, вы можете заставить Observable начать испускать элементы в нужное время.
Оператор RefCount автоматизирует процесс подключения и отключения от connectable Observable. Он работает с connectable Observable и возвращает обычное Observable. Когда первый наблюдатель подписывается на это Observable, RefCount подключается к базовому connectable Observable. RefCount отслеживает количество других наблюдателей, подписывающихся на него, и не отключается от базового connectable Observable до тех пор, пока последний наблюдатель не сделает этого.
См. также
- Connect
- Publish
- Replay
- Введение в Rx: RefCount
- Праздничный ужин: Share, Publish, Refcount и всё такое прочее от Каушика Гопала
Информация, специфичная для языка
RxGroovy refCount share
RxGroovy реализует этот оператор как refCount.
- Javadoc:
refCount()
Также есть оператор share, который эквивалентен применению операторов publish и refCount к Observable в указанном порядке.
- Javadoc:
share()
RxJava 1․x refCount share
RxJava реализует этот оператор как refCount.
- Javadoc:
refCount()
Также есть оператор share, который эквивалентен применению операторов publish и refCount к Observable в указанном порядке.
- Javadoc:
share()
RxJS refCount share shareValue
RxJava реализует этот оператор как refCount.
Пример кода
var interval = Rx.Observable.interval(1000);
var source = interval
.take(2)
.doAction(function (x) { console.log('Side effect'); });
var published = source.publish().refCount();
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 Next: SourceB0 Side effect Next: SourceA1 Next: SourceB1 Completed Completed
refCount находится в следующих дистрибутивах:
rx.all.jsrx.all.compat.js-
rx.binding.js(требуетrx.js,rx.compat.js,rx.lite.js, илиrx.lite.compat.js) rx.lite.jsrx.lite.compat.js
Также есть оператор share, который эквивалентен применению операторов publish и refCount к Observable в указанном порядке. Вариант под названием shareValue принимает в качестве параметра один элемент, который он испустит всем подписчикам перед началом испускания элементов из исходного Observable.
Пример кода
var interval = Rx.Observable.interval(1000);
var source = interval
.take(2)
.do(
function (x) { console.log('Side effect'); });
var published = source.share();
// When the number of observers subscribed to published observable goes from
// 0 to 1, we connect to the underlying observable sequence.
published.subscribe(createObserver('SourceA'));
// When the second subscriber is added, no additional subscriptions are added to the
// underlying observable sequence. As a result the operations that result in side
// effects are not repeated per subscriber.
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 Next: SourceB0 Side effect Next: SourceA1 Next: SourceB1 Completed
share и shareValue находятся в следующих дистрибутивах:
rx.all.jsrx.all.compat.js-
rx.binding.js(требуетrx.jsилиrx.compat.js) rx.lite.jsrx.lite.compat.js
RxPHP share singleInstance shareValue
RxPHP реализует этот оператор как share.
Возвращает последовательность observable, которая разделяет одну подписку на базовую последовательность. Этот оператор является специализацией publish, которая создает подписку, когда количество наблюдателей меняется с нуля на единицу, затем разделяет эту подписку со всеми последующими наблюдателями, пока количество наблюдателей не вернется к нулю, в этот момент подписка удаляется.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/share/share.php
//With Share
$source = \Rx\Observable::interval(1000)
->take(2)
->doOnNext(function ($x) {
echo "Side effect\n";
});
$published = $source->share();
$published->subscribe($createStdoutObserver('SourceA '));
$published->subscribe($createStdoutObserver('SourceB ')); Side effect SourceA Next value: 0 SourceB Next value: 0 Side effect SourceA Next value: 1 SourceB Next value: 1 SourceA Complete! SourceB Complete!
RxPHP также имеет оператор singleInstance.
Возвращает последовательность observable, которая разделяет одну подписку на базовую последовательность. Эта последовательность observable может быть повторно подписана, даже если все предыдущие подписки закончились. Этот оператор ведет себя как share() в RxJS 5
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/share/singleInstance.php
$interval = Rx\Observable::interval(1000);
$source = $interval
->take(2)
->do(function () {
echo 'Side effect', PHP_EOL;
});
$single = $source->singleInstance();
// two simultaneous subscriptions, lasting 2 seconds
$single->subscribe($createStdoutObserver('SourceA '));
$single->subscribe($createStdoutObserver('SourceB '));
\Rx\Observable::timer(5000)->subscribe(function () use ($single, &$createStdoutObserver) {
// resubscribe two times again, more than 5 seconds later,
// long after the original two subscriptions have ended
$single->subscribe($createStdoutObserver('SourceC '));
$single->subscribe($createStdoutObserver('SourceD '));
}); Side effect SourceA Next value: 0 SourceB Next value: 0 Side effect SourceA Next value: 1 SourceB Next value: 1 SourceA Complete! SourceB Complete! Side effect SourceC Next value: 0 SourceD Next value: 0 Side effect SourceC Next value: 1 SourceD Next value: 1 SourceC Complete! SourceD Complete!
RxPHP также имеет оператор shareValue.
Возвращает последовательность observable, которая разделяет одну подписку на базовую последовательность и начинается с начального значения. Этот оператор является специализацией publishValue, которая создает подписку, когда количество наблюдателей меняется с нуля на единицу, затем разделяет эту подписку со всеми последующими наблюдателями, пока количество наблюдателей не вернется к нулю, в этот момент подписка удаляется.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/share/shareValue.php
$source = \Rx\Observable::interval(1000)
->take(2)
->doOnNext(function ($x) {
echo "Side effect\n";
});
$published = $source->shareValue(42);
$published->subscribe($createStdoutObserver('SourceA '));
$published->subscribe($createStdoutObserver('SourceB ')); SourceA Next value: 42 SourceB Next value: 42 Side effect SourceA Next value: 0 SourceB Next value: 0 Side effect SourceA Next value: 1 SourceB Next value: 1 SourceA Complete! SourceB Complete!
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/refcount.html