Spec-Zone.ru › ReactiveX

Публикация

преобразовать обычное Observable в соединяемое Observable
Publish

Соединяемое Observable напоминает обычное Observable, за исключением того, что оно не начинает излучать элементы при подписке, а только когда к нему применяется оператор Connect. Таким образом, вы можете заставить Observable начать излучать элементы в выбранное вами время.

См. также

  • Connect
  • RefCount
  • Replay
  • Введение в Rx: Публикация и соединение
  • 101 Примеры Rx: Публикация — совместное использование подписки с несколькими наблюдателями
  • Свадебная вечеринка: Общий доступ, Публикация, Refcount и все такое от Каушика Гопала

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

RxGroovy publish

publish

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

  • Javadoc: publish()
publish

Также существует вариант, который принимает функцию в качестве параметра. Эта функция принимает излучаемый элемент из исходного Observable в качестве параметра и генерирует элемент, который будет излучен в его место результатом Observable.

  • Javadoc: publish(Func1)

RxJava 1․x publish

publish

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

  • Javadoc: publish()
publish

Также существует вариант, который принимает функцию в качестве параметра. Эта функция принимает в качестве параметра ConnectableObservable, который разделяет одну подписку на основную последовательность Observable. Эта функция генерирует и возвращает новую последовательность Observable.

  • Javadoc: publish(Func1)

RxJS let letBind multicast publish publishLast publishValue

publish

В 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

Оператор 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

Оператор 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.js
  • rx.all.compat.js
  • rx.binding.js (требует либо rx.js, либо rx.compat.js)
  • rx.lite.js
  • rx.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.js
  • rx.all.compat.js
  • rx.binding.js (требует либо rx.lite.js, либо rx.compat.js)
  • rx.lite.js
  • rx.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.js
  • rx.all.compat.js
  • rx.experimental.js

Он требует одного из следующих пакетов:

  • rx.js
  • rx.compat.js
  • rx.lite.js
  • rx.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

Spec-Zone.ru

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