Публикация
преобразовать обычное Observable в соединяемое Observable
Соединяемое Observable напоминает обычное Observable, за исключением того, что оно не начинает излучать элементы при подписке, а только когда к нему применяется оператор Connect. Таким образом, вы можете заставить Observable начать излучать элементы в выбранное вами время.
См. также
- Connect
- RefCount
- Replay
- Введение в Rx: Публикация и соединение
- 101 Примеры Rx: Публикация — совместное использование подписки с несколькими наблюдателями
- Свадебная вечеринка: Общий доступ, Публикация, Refcount и все такое от Каушика Гопала
Информация, специфичная для языка
RxGroovy publish
RxGroovy реализует этот оператор как publish.
- Javadoc:
publish()
Также существует вариант, который принимает функцию в качестве параметра. Эта функция принимает излучаемый элемент из исходного Observable в качестве параметра и генерирует элемент, который будет излучен в его место результатом Observable.
- Javadoc:
publish(Func1)
RxJava 1․x publish
RxJava реализует этот оператор как publish.
- Javadoc:
publish()
Также существует вариант, который принимает функцию в качестве параметра. Эта функция принимает в качестве параметра ConnectableObservable, который разделяет одну подписку на основную последовательность Observable. Эта функция генерирует и возвращает новую последовательность Observable.
- Javadoc:
publish(Func1)
RxJS let letBind multicast publish publishLast publishValue
В RxJS оператор publish принимает функцию в качестве параметра. Эта функция принимает излучаемый элемент из исходного Observable в качестве параметра и генерирует элемент, который будет излучен на его место возвращаемым ConnectableObservable.
Пример кода
var interval = Rx.Observable.interval(1000);
var source = interval
.take(2)
.doAction(function (x) {
console.log('Side effect');
});
var published = source.publish();
published.subscribe(createObserver('SourceA'));
published.subscribe(createObserver('SourceB'));
var connection = published.connect();
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
Оператор publishValue, помимо функции, описанной выше, принимает начальный элемент, который будет излучен результатом ConnectableObservable во время подключения перед излучением элементов из исходного Observable. Однако он не будет излучать этот начальный элемент наблюдателям, которые подключаются после подключения.
Пример кода
var interval = Rx.Observable.interval(1000);
var source = interval
.take(2)
.doAction(function (x) {
console.log('Side effect');
});
var published = source.publishValue(42);
published.subscribe(createObserver('SourceA'));
published.subscribe(createObserver('SourceB'));
var connection = published.connect();
function createObserver(tag) {
return Rx.Observer.create(
function (x) { console.log('Next: ' + tag + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); });
}
Next: SourceA42 Next: SourceB42 Side effect Next: SourceA0 Next: SourceB0 Side effect Next: SourceA1 Next: SourceB1 Completed Completed
Оператор publishLast аналогичен publish и принимает функцию с аналогичным поведением в качестве параметра. Он отличается от publish тем, что вместо применения этой функции к и излучения элемента для каждого элемента, излучаемого исходным Observable после подключения, он применяет эту функцию и излучает элемент только для последнего элемента, излученного исходным Observable, когда этот исходный Observable завершается нормально.
Пример кода
var interval = Rx.Observable.interval(1000);
var source = interval
.take(2)
.doAction(function (x) {
console.log('Side effect');
});
var published = source.publishLast();
published.subscribe(createObserver('SourceA'));
published.subscribe(createObserver('SourceB'));
var connection = published.connect();
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 Side effect Next: SourceA1 Completed Next: SourceB1 Completed
Указанные выше операторы доступны в следующих пакетах:
rx.all.jsrx.all.compat.js-
rx.binding.js(требует либоrx.js, либоrx.compat.js) rx.lite.jsrx.lite.compat.js
RxJS также имеет оператор multicast, который работает с обычным Observable, мультиплицирует это Observable с помощью указанного вами Subject, применяет трансформирующую функцию к каждому излучению и затем излучает эти преобразованные значения как собственную обычную последовательность Observable. Каждая подписка на это новое Observable будет запускать новую подписку на основное мультиплицируемое Observable.
Пример кода
var subject = new Rx.Subject();
var source = Rx.Observable.range(0, 3)
.multicast(subject);
var observer = Rx.Observer.create(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }
);
var subscription = source.subscribe(observer);
subject.subscribe(observer);
var connected = source.connect();
subscription.dispose();
Next: 0 Next: 0 Next: 1 Next: 1 Next: 2 Next: 2 Completed
Оператор multicast доступен в следующих пакетах:
rx.all.jsrx.all.compat.js-
rx.binding.js(требует либоrx.lite.js, либоrx.compat.js) rx.lite.jsrx.lite.compat.js
Также существует оператор let (псевдоним letBind доступен для браузеров, таких как Internet Explorer до IE9, где «let» запрещен). Он похож на multicast, но не мультиплицирует основное Observable через Subject:
Пример кода
var obs = Rx.Observable.range(1, 3);
var source = obs.let(function (o) { return o.concat(o); });
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); });
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); });
Next: 1 Next: 2 Next: 3 Next: 1 Next: 2 Next: 3 Completed
Оператор let (или letBind) доступен в следующих пакетах:
rx.all.jsrx.all.compat.jsrx.experimental.js
Он требует одного из следующих пакетов:
rx.jsrx.compat.jsrx.lite.jsrx.lite.compat.js
RxPHP multicast multicastWithSelector publish publishLast publishValue
RxPHP реализует этот оператор как multicast.
Мультиплицирует уведомления о последовательности источника через созданный Subject ко всем использованиям последовательности внутри функции селектора. Каждая подписка на результирующую последовательность вызывает отдельное вызов мультипликации, экспонируя последовательность, полученную от вызова функции селектора. Для специализаций с фиксированными типами Subject см. Publish, PublishLast и Replay.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/multicast/multicast.php $subject = new \Rx\Subject\Subject(); $source = \Rx\Observable::range(0, 3)->multicast($subject); $subscription = $source->subscribe($stdoutObserver); $subject->subscribe($stdoutObserver); $connected = $source->connect();
Next value: 0 Next value: 0 Next value: 1 Next value: 1 Next value: 2 Next value: 2 Complete!
RxPHP также имеет оператор multicastWithSelector.
Мультиплицирует уведомления о последовательности источника через Subject, созданный с помощью фабрики селектора Subject, ко всем использованиям последовательности внутри функции селектора. Каждая подписка на результирующую последовательность вызывает отдельное вызов мультипликации, экспонируя последовательность, полученную от вызова функции селектора. Для специализаций с фиксированными типами Subject см. Publish, PublishLast и Replay.
RxPHP также имеет оператор publish.
Возвращает последовательность Observable, которая является результатом вызова селектора для соединяемой последовательности Observable, которая разделяет одну подписку на основную последовательность. Этот оператор является специализацией Multicast с использованием обычного Subject.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/publish/publish.php
/* With publish */
$interval = \Rx\Observable::range(0, 10);
$source = $interval
->take(2)
->doOnNext(function ($x) {
echo "Side effect\n";
});
$published = $source->publish();
$published->subscribe($createStdoutObserver('SourceC '));
$published->subscribe($createStdoutObserver('SourceD '));
$published->connect();
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 также имеет оператор publishLast.
Возвращает последовательность Observable, которая является результатом вызова селектора для соединяемой последовательности Observable, которая разделяет одну подписку на основную последовательность, содержащую только последнее уведомление. Этот оператор является специализацией Multicast с использованием AsyncSubject.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/publish/publishLast.php
$range = \Rx\Observable::fromArray(range(0, 1000));
$source = $range
->take(2)
->doOnNext(function ($x) {
echo "Side effect\n";
});
$published = $source->publishLast();
$published->subscribe($createStdoutObserver('SourceA'));
$published->subscribe($createStdoutObserver('SourceB'));
$connection = $published->connect();
Side effect Side effect SourceANext value: 1 SourceBNext value: 1 SourceAComplete! SourceBComplete!
RxPHP также имеет оператор publishValue.
Возвращает последовательность Observable, которая является результатом вызова селектора для соединяемой последовательности Observable, которая разделяет одну подписку на основную последовательность и начинается с initialValue. Этот оператор является специализацией Multicast с использованием BehaviorSubject.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/publish/publishValue.php
$range = \Rx\Observable::fromArray(range(0, 1000));
$source = $range
->take(2)
->doOnNext(function ($x) {
echo "Side effect\n";
});
$published = $source->publishValue(42);
$published->subscribe($createStdoutObserver('SourceA'));
$published->subscribe($createStdoutObserver('SourceB'));
$connection = $published->connect();
SourceANext value: 42 SourceBNext value: 42 Side effect SourceANext value: 0 SourceBNext value: 0 Side effect SourceANext value: 1 SourceBNext value: 1 SourceAComplete! SourceBComplete!
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/publish.html