Объединение
объединяет несколько Observable в один, объединяя их эмиссии
Вы можете объединить выходные данные нескольких Observable, чтобы они действовали как одно Observable, используя оператор Merge.
Merge может чередовать элементы, испускаемые объединёнными Observable (аналогичный оператор, Concat, не чередует элементы, но испускает все элементы каждого исходного Observable по очереди, прежде чем начать испускать элементы из следующего исходного Observable).
Как показано на диаграмме выше, уведомление onError от любого из исходных Observable немедленно будет передано наблюдателям и завершит объединённое Observable.
Во многих реализациях ReactiveX существует второй оператор, MergeDelayError, который изменяет это поведение — сохраняя onError уведомления до завершения всех объединённых Observable и только тогда передаёт их наблюдателям:
См. также
Информация, специфичная для языка
RxClojure interleave interleave* merge merge* merge-delay-error merge-delay-error*
В RxClojure здесь имеется шесть операторов, которые стоит рассмотреть:
merge преобразует два или более Observable в одно Observable, которое испускает все элементы, испускаемые всеми этими Observable.
merge* преобразует Observable, испускающее Observable, в одно Observable, которое испускает все элементы, испускаемые всеми испущенными Observable.
merge-delay-error подобен merge, но будет испускать все элементы всех объединённых Observable, даже если одно или несколько из этих Observable завершаются уведомлением onError при наличии ожидаемых эмиссий.
merge-delay-error* — это аналогично изменённая версия merge*.
interleave подобен merge, но более продуманно чередует элементы из исходных Observable: результирующее Observable испускает первый элемент, испущенный первым исходным Observable, затем первый элемент, испущенный вторым исходным Observable, и так далее, и дойдя до последнего исходного Observable, затем испускает второй элемент, испущенный первым исходным Observable, второй элемент, испущенный вторым исходным Observable, и так далее, пока все исходные Observable не завершатся.
interleave* аналогичен, но работает с Observable из Observable.
RxCpp merge
RxCpp реализует этот оператор как merge.
RxGroovy merge mergeDelayError mergeWith
RxGroovy реализует этот оператор как merge, mergeWith, и mergeDelayError.
Например, следующий код объединяет odds и evens в одно Observable. (Оператор subscribeOn заставляет odds работать на другом потоке из evens, чтобы обе Observable могли испускать элементы одновременно, чтобы продемонстрировать, как Merge может чередовать эти элементы.)
Пример кода
odds = Observable.from([1, 3, 5, 7]).subscribeOn(someScheduler);
evens = Observable.from([2, 4, 6]);
Observable.merge(odds,evens).subscribe(
{ println(it); }, // onNext
{ println("Error: " + it.getMessage()); }, // onError
{ println("Sequence complete"); } // onCompleted
); 1 3 2 5 4 7 6 Sequence complete
Вместо передачи нескольких Observable (до девяти) в merge, вы также можете передать List<> (или другой Iterable) Observable, массив Observable или даже Observable, испускающее Observable, и merge объединит их выходные данные в выходные данные одного Observable:
- Javadoc:
merge(Iterable) - Javadoc:
merge(Iterable,int) - Javadoc:
merge(Observable[]) - Javadoc:
merge(Observable[], int)(RxGroovy 1.1) - Javadoc:
merge(Observable, Observable)(есть также версии, принимающие до девяти Observable)
Если вы передаёте Observable из Observable, у вас есть возможность также передать значение, указывающее merge максимальное количество этих Observable, к которым он должен пытаться подписаться одновременно. Достигнув этого максимального числа подписок, он воздержится от подписки на любые другие Observable, испущенные исходным Observable, до тех пор, пока одно из уже подписанных Observable не выпустит уведомление onCompleted.
Версия оператора merge — mergeWith, так что, например, в примере кода выше вместо Observable.merge(odds,evens) можно написать odds.mergeWith(evens).
- Javadoc:
mergeWith(Observable)
Если какое-либо из отдельных Observable завершится уведомлением onError, Observable, созданное merge, немедленно завершится уведомлением onError. Если вы предпочитаете объединение, которое продолжает испускать результаты оставшихся, не содержащих ошибок Observable, прежде чем сообщить об ошибке, используйте mergeDelayError вместо этого.
mergeDelayError ведет себя очень похоже на merge. Исключение составляет ситуация, когда одно из объединяемых Observable завершается уведомлением onError. Если это произойдёт с merge, объединённое Observable немедленно выпустит уведомление onError и завершится. mergeDelayError, с другой стороны, отложит сообщение об ошибке, пока не даст другим не генерирующим ошибки Observable, которые он объединяет, шанс завершить испускание их элементов, и он выпустит их сам, и завершится уведомлением onError только тогда, когда все другие объединённые Observable завершатся.
Поскольку возможно, что более чем одно из объединённых Observable столкнулось с ошибкой, mergeDelayError может передать информацию о нескольких ошибках в уведомлении onError (он никогда не вызовет метод наблюдателя onError более одного раза). По этой причине, если вы хотите узнать о природе этих ошибок, вы должны написать методы наблюдателей onError таким образом, чтобы они принимали параметр класса CompositeException.
mergeDelayError имеет меньше вариантов. Вы не можете передать ему Iterable или массив Observable, но вы можете передать ему Observable, испускающее Observable, или от одного до девяти отдельных Observable в качестве параметров. Нет версии метода экземпляра для mergeDelayError, как для merge.
- Javadoc:
mergeDelayError(Observable<Observable>) - Javadoc:
mergeDelayError(Observable,Observable)(есть также версии, принимающие до девяти Observable)
RxJava 1․x merge mergeDelayError mergeWith
RxJava реализует этот оператор как merge, mergeWith, и mergeDelayError.
Пример кода
Observable<Integer> odds = Observable.just(1, 3, 5).subscribeOn(someScheduler);
Observable<Integer> evens = Observable.just(2, 4, 6);
Observable.merge(odds, evens)
.subscribe(new Subscriber<Integer>() {
@Override
public void onNext(Integer item) {
System.out.println("Next: " + item);
}
@Override
public void onError(Throwable error) {
System.err.println("Error: " + error.getMessage());
}
@Override
public void onCompleted() {
System.out.println("Sequence complete.");
}
}); Next: 1 Next: 3 Next: 5 Next: 2 Next: 4 Next: 6 Sequence complete.
- Javadoc:
merge(Iterable) - Javadoc:
merge(Iterable,int) - Javadoc:
merge(Observable[]) - Javadoc:
merge(Observable[], int)(RxJava 1.1) - Javadoc:
merge(Observable, Observable)(также существуют версии, принимающие до девяти Observables)
Вместо передачи нескольких Observables (до девяти) в merge, вы также можете передать List<> (или другой Iterable) Observables, массив Observables, или даже Observable, испускающий Observables, и merge объединит их выходные данные в выходные данные одного Observable:
Если вы передаёте Observable из Observables, у вас есть возможность также передать значение, указывающее merge максимальное количество таких Observables, к которым он должен попытаться подписаться одновременно. После достижения этого максимального количества подписок, он откажется от подписки на любые другие Observables, испускаемые исходным Observable, до тех пор, пока одно из уже подписанных Observables не выдаст onCompleted уведомление.
- Javadoc:
merge(Observable<Observable>) - Javadoc:
merge(Observable<Observable>, int)(RxJava 1.1)
Версия оператора merge это mergeWith, поэтому, например, вместо записи Observable.merge(odds,evens) вы также можете написать odds.mergeWith(evens).
Если любой из отдельных Observables, переданных в merge, завершается onError уведомлением, Observable, производимое merge, немедленно завершится onError уведомлением. Если вам требуется объединение, которое продолжает испускать результаты оставшихся без ошибок Observables перед сообщением об ошибке, используйте mergeDelayError вместо этого.
mergeDelayError ведет себя очень похоже на merge. Единственное исключение - когда одно из объединяемых Observables завершается onError уведомлением. Если это произойдёт с merge, объединённое Observable немедленно выдаст onError уведомление и завершится. mergeDelayError, с другой стороны, отложит сообщение об ошибке, пока не даст всем другим Observables, которые он объединяет, возможность завершить испускание своих элементов, и он выпустит их сам, и завершится onError уведомлением только тогда, когда все другие объединённые Observables завершат работу.
Поскольку возможно, что более чем одно из объединённых Observables столкнулось с ошибкой, mergeDelayError может передать информацию о нескольких ошибках в onError уведомлении (он никогда не вызовет метод наблюдателя onError более одного раза). По этой причине, если вы хотите узнать о природе этих ошибок, вы должны написать методы наблюдателей onError таким образом, чтобы они принимали параметр класса CompositeException.
mergeDelayError имеет меньше вариантов. Вы не можете передать ему Iterable или массив Observables, но вы можете передать ему Observable, который испускает Observables, или от одного до девяти отдельных Observables в качестве параметров. Нет версии метода экземпляра для mergeDelayError, как для merge.
- Javadoc:
mergeDelayError(Observable<Observable>) - Javadoc:
mergeDelayError(Observable,Observable)(также существуют версии, принимающие до девяти Observables)
RxJS merge mergeAll mergeDelayError
Первый вариант merge - это оператор экземпляра, который принимает переменное число Observables в качестве параметров, объединяя каждый из этих Observables с исходными (экземпляра) Observables для создания одного выходного Observable.
Этот первый вариант merge встречается в следующих дистрибутивах:
rx.jsrx.compat.jsrx.lite.jsrx.lite.compat.js
Второй вариант merge - это оператор прототипа (класса), который принимает два параметра. Второй из них - Observable, который испускает Observables, которые вы хотите объединить. Первый - число, указывающее максимальное количество этих испускаемых Observables, которые вы хотите merge попытаться подписаться на них в любой момент. После достижения этого максимального количества подписок, он откажется от подписки на любые другие Observables, испускаемые исходным Observable, до тех пор, пока одно из уже подписанных Observables не выдаст onCompleted уведомление.
Этот второй вариант merge встречается в следующих дистрибутивах:
rx.jsrx.all.jsrx.all.compat.jsrx.compat.jsrx.lite.jsrx.lite.compat.js
mergeAll похож на этот второй вариант merge, за исключением того, что не позволяет установить это максимальное количество подписок. Он принимает только один параметр - Observable из Observables.
mergeAll встречается в следующих дистрибутивах:
rx.jsrx.all.jsrx.all.compat.jsrx.compat.jsrx.lite.jsrx.lite.compat.js
Если любой из отдельных Observables, переданных в merge или mergeAll, завершается onError уведомлением, результирующее Observable немедленно завершится onError уведомлением. Если вам требуется объединение, которое продолжает испускать результаты оставшихся, без ошибок Observables перед сообщением об ошибке, используйте mergeDelayError вместо этого.
Пример кода
var source1 = Rx.Observable.of(1,2,3);
var source2 = Rx.Observable.throwError(new Error('whoops!'));
var source3 = Rx.Observable.of(4,5,6);
var merged = Rx.Observable.mergeDelayError(source1, source2, source3);
var subscription = merged.subscribe(
function (x) { console.log('Next: %s', x); },
function (err) { console.log('Error: %s', err); }
function () { console.log('Completed' } ); 1 2 3 4 5 6 Error: Error: whoops!
mergeDelayError встречается в следующих дистрибутивах:
rx.jsrx.compat.jsrx.lite.jsrx.lite.compat.js
RxKotlin merge mergeDelayError mergeWith
RxKotlin реализует этот оператор как merge, mergeWith, и mergeDelayError.
Вместо передачи нескольких Observables (до девяти) в merge, вы также можете передать List<> (или другой Iterable) Observables, массив Observables, или даже Observable, которое испускает Observables, и merge будет объединять их выходные данные в выходные данные одного Observable:
Если вы передаёте Observable Observables, у вас есть возможность передать значение, указывающее merge максимальное количество этих Observables, к которым оно должно попытаться подключиться одновременно. Достигнув этого максимального количества подписок, оно воздержится от подписки на другие Observables, испускаемые исходным Observable, до тех пор, пока одно из уже подписанных Observables не выдаст onCompleted уведомление.
Вариант оператора merge — mergeWith, поэтому, например, вместо записи Observable.merge(odds,evens) вы также можете написать odds.mergeWith(evens).
Если любое из отдельных Observables, переданных в merge, завершится уведомлением об onError, Observable, созданный merge, немедленно завершится уведомлением об onError. Если вы предпочтёте объединение, которое продолжает испускать результаты оставшихся, не содержащих ошибок Observables, прежде чем сообщать об ошибке, используйте mergeDelayError вместо этого.
mergeDelayError ведет себя очень похоже на merge. Разница заключается в том, что если одно из объединяемых Observables завершается уведомлением об onError . Если это происходит с merge, объединённое Observable немедленно выдаёт уведомление об onError и завершается. mergeDelayError, с другой стороны, отложит сообщение об ошибке, пока не предоставит другим, не генерирующим ошибок, Observable, которые оно объединяет, возможность закончить испускание своих элементов, и само их испустит, и завершится уведомлением об onError только тогда, когда все остальные объединяемые Observables закончат свою работу.
Поскольку возможно, что более чем одно из объединённых Observables столкнулось с ошибкой, mergeDelayError может передать информацию о нескольких ошибках в уведомлении об onError (он никогда не вызовет метод наблюдателя onError более одного раза). Поэтому, если вы хотите узнать о природе этих ошибок, вам следует написать методы наблюдателей onError таким образом, чтобы они принимали параметр класса CompositeException.
mergeDelayError имеет меньше вариантов. Вы не можете передать ему Iterable или массив Observables, но можете передать Observable, испускающий Observables, или от одного до девяти отдельных Observables в качестве параметров. Нет версии метода экземпляра для mergeDelayError как для merge.
RxNET Merge
Rx.NET реализует этот оператор как Merge.
Вы можете передать Merge массив Observables, перечислитель Observables, Observable Observables или две отдельные Observables.
Если вы передаёте перечислитель или Observable Observables, у вас есть возможность передать целое число, указывающее максимальное количество этих Observables, к которым оно должно попытаться подключиться одновременно. Достигнув этого максимального количества подписок, оно воздержится от подписки на другие Observables, испускаемые исходным Observable, до тех пор, пока одно из уже подписанных Observables не выдаст onCompleted уведомление.
RxPHP merge mergeAll
RxPHP реализует этот оператор как merge.
Объедините Observable с другим Observable, объединив их испускания в одно Observable.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/merge/merge.php
$observable = Rx\Observable::of(42)->repeat();
$otherObservable = Rx\Observable::of(21)->repeat();
$mergedObservable = $observable
->merge($otherObservable)
->take(10);
$disposable = $mergedObservable->subscribe($stdoutObserver); Next value: 42 Next value: 21 Next value: 42 Next value: 21 Next value: 42 Next value: 21 Next value: 42 Next value: 21 Next value: 42 Next value: 21 Complete!
RxPHP также имеет оператор mergeAll.
Объединяет последовательность Observable в Observable.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/merge/merge-all.php
$sources = Rx\Observable::range(0, 3)
->map(function ($x) {
return Rx\Observable::range($x, 3);
});
$merged = $sources->mergeAll();
$disposable = $merged->subscribe($stdoutObserver); Next value: 0 Next value: 1 Next value: 1 Next value: 2 Next value: 2 Next value: 2 Next value: 3 Next value: 3 Next value: 4 Complete!
RxPY merge merge_all merge_observable
RxPY реализует этот оператор как merge и merge_all/merge_observable.
Вы можете либо передать merge набор Observables в качестве отдельных параметров, либо в качестве единственного параметра, содержащего массив этих Observables.
merge_all и его псевдоним merge_observable принимают в качестве единственного параметра Observable, который испускает Observables. Они объединяют испускания всех этих Observables для создания собственного Observable.
Rxrb merge merge_all merge_concurrent
Rx.rb реализует этот оператор как merge, merge_concurrent, и merge_all.
merge объединяет второе Observable в то, которое обрабатывается, чтобы создать новое объединённое Observable.
merge_concurrent работает с Observable, который испускает Observables, объединяя испускания каждого из этих Observables в собственные испускания. Вы можете необязательно передать ему целочисленный параметр, указывающий, сколько из этих испущенных Observables merge_concurrent должно пытаться подписаться одновременно. Достигнув этого максимального числа подписок, оно воздержится от подписки на другие Observables, испускаемые исходным Observable, до тех пор, пока одно из уже подписанных Observables не выдаст onCompleted уведомление. По умолчанию это 1, что делает его эквивалентным merge_all.
merge_all похож на merge_concurrent(1). Он подписывается на каждое испущенное Observable по одному, отображая его испускания как свои собственные и ожидая подписки на следующее Observable до тех пор, пока текущее не завершится уведомлением об onCompleted . В этом отношении он больше похож на вариант Concat.
RxScala flatten flattenDelayError merge mergeDelayError
RxScala реализует этот оператор как flatten, flattenDelayError, merge, и mergeDelayError.
merge принимает второе Observable в качестве параметра и объединяет это Observable с тем, к которому применяется оператор merge, чтобы создать новое выходное Observable.
mergeDelayError похож на merge, за исключением того, что он всегда испускает все элементы из обоих Observables, даже если одно из Observables завершается уведомлением об onError до того, как другое Observable закончит испускание элементов.
flatten принимает в качестве параметра Observable, который испускает Observables. Он объединяет элементы, испускаемые каждым из этих Observables, чтобы создать собственную последовательность единого Observable. Вариант этого оператора позволяет вам передать Int указывающее максимальное количество из этих испущенных Observables, которые вы хотите flatten попробовать подписаться одновременно. Если он достигнет этого максимального количества подписок, он воздержится от подписки на другие Observables, испускаемые исходным Observable, до тех пор, пока одно из уже подписанных Observables не выдаст onCompleted уведомление.
flattenDelayError похож на flatten за исключением того, что он всегда испускает все элементы из всех испущенных Observables, даже если одно или несколько из этих Observables завершаются уведомлением об onError до того, как другие Observables закончат испускание элементов.
RxSwift merge
RxSwift реализует этот оператор как merge.
merge принимает в качестве параметра Observable, который испускает Observables. Он объединяет элементы, испускаемые каждым из этих Observables, для создания собственной последовательности единого Observable.
Вариант этого оператора merge(maxConcurrent:) позволяет вам передать Int указывающее максимальное количество из этих испущенных Observables, которые вы хотите merge попробовать подписаться одновременно. При достижении этого максимального числа подписок, он воздержится от подписки на другие Observables, испускаемые исходным Observable, до тех пор, пока одно из уже подписанных Observables не выдаст onCompleted уведомление.
Пример кода
let subject1 = PublishSubject()
let subject2 = PublishSubject()
Observable.of(subject1, subject2)
.merge()
.subscribe {
print($0)
}
subject1.on(.Next(10))
subject1.on(.Next(11))
subject1.on(.Next(12))
subject2.on(.Next(20))
subject2.on(.Next(21))
subject1.on(.Next(14))
subject1.on(.Completed)
subject2.on(.Next(22))
subject2.on(.Completed) Next(10) Next(11) Next(12) Next(20) Next(21) Next(14) Next(22) Completed
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/merge.html