Zip
объединяет эмиссии нескольких Observables вместе с помощью указанной функции и выпускает одиночные элементы для каждого сочетания на основе результатов этой функции
Метод Zip возвращает Observable, который применяет функцию по вашему выбору к сочетанию элементов, выпущенных последовательно, двумя (или более) другими Observables, при этом результаты этой функции становятся элементами, выпускаемыми возвращаемым Observable. Он применяет эту функцию в строгой последовательности, поэтому первый элемент, выпущенный новым Observable, будет результатом функции, применённой к первому элементу, выпущенному Observable #1, и первому элементу, выпущенному Observable #2; второй элемент, выпущенный новым zip-Observable, будет результатом функции, применённой ко второму элементу, выпущенному Observable #1, и второму элементу, выпущенному Observable #2; и так далее. Он будет выпускать только столько элементов, сколько элементов выпущено источником Observable, который выпускает наименьшее количество элементов.
См. также
Информация, специфичная для языка
RxGroovy zip zipWith
RxGroovy реализует этот оператор как несколько вариантов zip и также как zipWith, версия оператора в виде метода экземпляра.
Последний аргумент к zip — это функция, которая принимает элемент из каждого из Observables, которые объединяются и выпускает элемент, который должен быть выпущен в ответ Observable, возвращаемым из zip. Вы можете предоставить Observables для объединения в zip либо как от двух до девяти отдельных параметров, либо как один параметр: либо Iterable из Observables, либо Observable, который выпускает Observables (как на иллюстрации выше).
Пример кода
odds = Observable.from([1, 3, 5, 7, 9]);
evens = Observable.from([2, 4, 6]);
Observable.zip(odds, evens, {o, e -> [o, e]}).subscribe(
{ println(it); }, // onNext
{ println("Error: " + it.getMessage()); }, // onError
{ println("Sequence complete"); } // onCompleted
); [1, 2] [3, 4] [5, 6] Sequence complete
Обратите внимание, что в этом примере возвращаемое Observable завершается нормально после выдачи трёх элементов, что соответствует количеству элементов, выпущенных более коротким из двух исходных Observables (evens, который выпускает три элемента).
- Javadoc:
zip(Iterable<Observable>,FuncN) - Javadoc:
zip(Observable<Observable>,FuncN) - Javadoc:
zip(Observable,Observable,Func2)(есть также версии, принимающие до девяти Observables)
Версия оператора в виде метода экземпляра zipWith всегда принимает два параметра. Первый параметр может быть как простым Observable, так и Iterable (как на иллюстрации выше).
- Javadoc:
zipWith(Observable,Func2) - Javadoc:
zipWith(Iterable,Func2)
zip и zipWith по умолчанию не работают с каким-либо конкретным Scheduler.
RxJava 1․x zip zipWith
RxJava реализует этот оператор как несколько вариантов zip и также как zipWith, версию оператора в виде метода экземпляра.
Последний аргумент к zip — это функция, которая принимает элемент из каждого из Observables, которые объединяются, и выпускает элемент, который должен быть выпущен в ответ Observable, возвращаемым из zip. Вы можете предоставить Observables для объединения в zip либо как от двух до девяти отдельных параметров, либо как один параметр: либо Iterable из Observables, либо Observable, который выпускает Observables (как на иллюстрации выше).
- Javadoc:
zip(Iterable<Observable>,FuncN) - Javadoc:
zip(Observable<Observable>,FuncN) - Javadoc:
zip(Observable,Observable,Func2)(есть также версии, принимающие до девяти Observables)
Версия оператора в виде метода экземпляра zipWith всегда принимает два параметра. Первый параметр может быть как простым Observable, так и Iterable (как на иллюстрации выше).
- Javadoc:
zipWith(Observable,Func2) - Javadoc:
zipWith(Iterable,Func2)
zip и zipWith по умолчанию не работают с каким-либо конкретным Scheduler.
RxJS forkJoin zip zipArray
RxJS реализует этот оператор как zip и zipArray.
zip принимает переменное количество Observables или Promises в качестве параметров, за которым следует функция, принимающая один элемент, выпущенный каждым из этих Observables или разрешенный этими Promises в качестве входных данных, и генерирующая один элемент, который должен быть выпущен возвращаемым Observable.
Пример кода
/* Using arguments */
var range = Rx.Observable.range(0, 5);
var source = Observable.zip(
range,
range.skip(1),
range.skip(2),
function (s1, s2, s3) {
return s1 + ':' + s2 + ':' + s3;
}
);
var subscription = source.subscribe(
function (x) {
console.log('Next: ' + x);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}); Next: 0:1:2 Next: 1:2:3 Next: 2:3:4 Completed
/* Using promises and Observables */
var range = Rx.Observable.range(0, 5);
var source = Observable.zip(
RSVP.Promise.resolve(0),
RSVP.Promise.resolve(1),
Rx.Observable.return(2)
function (s1, s2, s3) {
return s1 + ':' + s2 + ':' + s3;
}
);
var subscription = source.subscribe(
function (x) {
console.log('Next: ' + x);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}); Next: 0:1:2 Completed
zipArray принимает переменное количество Observables в качестве параметров и возвращает Observable, который выпускает массивы, каждый из которых содержит n-й элемент из каждого исходного Observable.
Пример кода
var range = Rx.Observable.range(0, 5);
var source = Rx.Observable.zipArray(
range,
range.skip(1),
range.skip(2)
);
var subscription = source.subscribe(
function (x) {
console.log('Next: ' + x);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
}); Next: [0,1,2] Next: [1,2,3] Next: [2,3,4] Completed
RxJS также реализует аналогичный оператор forkJoin. Существуют две разновидности этого оператора. Первый собирает последний элемент, выпущенный каждым из исходных Observables, в массив и выпускает этот массив как свой единственный выпущенный элемент. Вы можете передать список Observables в forkJoin либо в качестве отдельных параметров, либо как массив Observables.
var source = Rx.Observable.forkJoin(
Rx.Observable.return(42),
Rx.Observable.range(0, 10),
Rx.Observable.fromArray([1,2,3]),
RSVP.Promise.resolve(56)
);
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }); Next: [42, 9, 3, 56] Completed
Существует вторая разновидность оператора forkJoin, представленная в виде прототипа функции, и вы вызываете её для экземпляра одного исходного Observable, передавая другой исходный Observable в качестве параметра. В качестве второго параметра вы передаёте функцию, которая объединяет последний элемент, выпущенный двумя исходными Observables, в единственный элемент, который должен быть выпущен возвращаемым Observable.
var source1 = Rx.Observable.return(42);
var source2 = Rx.Observable.range(0, 3);
var source = source1.forkJoin(source2, function (s1, s2) {
return s1 + s2;
});
var subscription = source.subscribe(
function (x) { console.log('Next: ' + x); },
function (err) { console.log('Error: ' + err); },
function () { console.log('Completed'); }); Next: 44 Completed
forkJoin присутствует в следующих распределениях:
rx.all.jsrx.all.compat.js-
rx.experimental.js(требуетrx.js,rx.compat.js,rx.lite.js, илиrx.lite.compat.js)
RxPHP zip forkJoin
RxPHP реализует этот оператор как zip.
Объединяет указанные последовательности Observable в одну последовательность Observable, используя функцию селектора всякий раз, когда все последовательности Observable произвели элемент в соответствующем индексе. Если функция селектора результата опущена, будет выпущен список с элементами последовательностей Observable в соответствующих индексах.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/zip/zip.php
//Without a result selector
$range = \Rx\Observable::fromArray(range(0, 4));
$source = $range
->zip([
$range->skip(1),
$range->skip(2)
]);
$observer = $createStdoutObserver();
$subscription = $source
->subscribe(new CallbackObserver(
function ($array) use ($observer) {
$observer->onNext(json_encode($array));
},
[$observer, 'onError'],
[$observer, 'onCompleted']
)); Next value: [0,1,2] Next value: [1,2,3] Next value: [2,3,4] Complete!
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/zip/zip-result-selector.php
//With a result selector
$range = \Rx\Observable::fromArray(range(0, 4));
$source = $range
->zip([
$range->skip(1),
$range->skip(2)
], function ($s1, $s2, $s3) {
return $s1 . ':' . $s2 . ':' . $s3;
});
$observer = $createStdoutObserver();
$subscription = $source->subscribe($createStdoutObserver()); Next value: 0:1:2 Next value: 1:2:3 Next value: 2:3:4 Complete!
RxPHP также имеет оператор forkJoin.
Выполняет все последовательности Observable параллельно и собирает их последние элементы.
Пример кода
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/forkJoin/forkJoin.php
use Rx\Observable;
$obs1 = Observable::range(1, 4);
$obs2 = Observable::range(3, 5);
$obs3 = Observable::fromArray(['a', 'b', 'c']);
$observable = Observable::forkJoin([$obs1, $obs2, $obs3], function($v1, $v2, $v3) {
return $v1 . $v2 . $v3;
});
$observable->subscribe($stdoutObserver); Next value: 47c Complete!
© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/zip.html