Spec-Zone.ru › ReactiveX

GroupBy

разделить Observable на набор Observables, каждый из которых излучает подмножество элементов из исходного Observable
GroupBy

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

См. также

  • Window
  • Введение в Rx: GroupBy
  • Анимации операторов Rx: GroupBy от ТамИра Дрешера

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

RxGroovy groupBy

groupBy

RxGroovy реализует оператор groupBy. Observable, которое он возвращает, излучает элементы определённого подкласса Observable — GroupedObservable. Объекты, реализующие интерфейс GroupedObservable, имеют дополнительный метод — getkey, с помощью которого можно получить ключ, по которому элементы были назначены этому конкретному GroupedObservable.

Следующий фрагмент кода использует groupBy для преобразования списка чисел в два списка, сгруппированных по тому, являются ли числа чётными или нет:

Пример кода

def numbers = Observable.from([1, 2, 3, 4, 5, 6, 7, 8, 9]);
def groupFunc = { return(0 == (it % 2)); };

numbers.groupBy(groupFunc).flatMap({ it.reduce([it.getKey()], {a, b -> a << b}) }).subscribe(
  { println(it); },                          // onNext
  { println("Error: " + it.getMessage()); }, // onError
  { println("Sequence complete"); }          // onCompleted
);
[false, 1, 3, 5, 7, 9]
[true, 2, 4, 6, 8]
Sequence complete

Другая версия groupBy позволяет передать функцию преобразования, которая изменяет элементы до их излучения результативными GroupedObservable.

Обратите внимание, что когда groupBy разбивает исходный Observable на Observable, излучающий GroupedObservables, каждый из этих GroupedObservables начинает буферизовать элементы, которые он будет излучать при подписке. По этой причине, если вы игнорируете какой-либо из этих GroupedObservables (вы не подписываетесь на него и не применяете к нему оператор, который подписывается на него), это может привести к потенциальной утечке памяти. По этой причине, вместо игнорирования GroupedObservables, на который вы не хотите подписаться, вы должны применить оператор, например, take(0), чтобы сигнализировать ему, что он может отбросить свой буфер.

Если вы отписываетесь от одного из GroupedObservables или если оператор, такой как takes, который вы применяете к GroupedObservables, отписывается от него, этот GroupedObservable будет завершён. Если исходный Observable позже излучит элемент, ключ которого соответствует GroupedObservables, который был завершён таким образом, groupBy создаст и излучит новый GroupedObservable для соответствия ключу. Другими словами, отписка от GroupedObservable не заставит groupBy пропустить элементы из его группы. Например, см. следующий код:

Пример кода

Observable.range(1,5)
          .groupBy({ 0 })
          .flatMap({ this.take(1) })
          .subscribe(
  { println(it); },                          // onNext
  { println("Error: " + it.getMessage()); }, // onError
  { println("Sequence complete"); }          // onCompleted
);
1
2
3
4
5

В приведенном выше коде исходный Observable излучает последовательность { 1 2 3 4 5 }. Когда он излучает первый элемент в этой последовательности, оператор groupBy создаёт и излучает GroupedObservable с ключом 0. Оператор flatMap применяет оператор take(1) к этому GroupedObservable, что даёт ему элемент (1), который он излучает и который также отписывается от GroupedObservable, который завершается. Когда исходный Observable излучает второй элемент в своей последовательности, оператор groupBy создаёт и излучает второй GroupedObservable с тем же ключом (0), чтобы заменить тот, который был завершён. flatMap снова применяет take(1) к этому новому GroupedObservable, чтобы получить новый элемент для излучения (2) и отписаться от и завершить GroupedObservable, и этот процесс повторяется для оставшихся элементов в исходной последовательности.

groupBy по умолчанию не работает с каким-либо конкретным Scheduler.

  • Javadoc: groupBy(Func1)
  • Javadoc: groupBy(Func1,Func1)

RxJava 1․x groupBy

groupBy

RxJava реализует оператор groupBy. Observable, которое он возвращает, излучает элементы определённого подкласса Observable — GroupedObservable. Объекты, реализующие интерфейс GroupedObservable, имеют дополнительный метод — getkey, с помощью которого можно получить ключ, по которому элементы были назначены этому конкретному GroupedObservable.

Другая версия groupBy позволяет передать функцию преобразования, которая изменяет элементы до их излучения результативными GroupedObservables.

Обратите внимание, что когда groupBy разбивает исходный Observable на Observable, излучающий GroupedObservables, каждый из этих GroupedObservables начинает буферизовать элементы, которые он будет излучать при подписке. По этой причине, если вы игнорируете какой-либо из этих GroupedObservables (вы не подписываетесь на него и не применяете к нему оператор, который подписывается на него), это может привести к потенциальной утечке памяти. По этой причине, вместо игнорирования GroupedObservables, на который вы не хотите подписаться, вы должны применить оператор, например, take(0), чтобы сигнализировать ему, что он может отбросить свой буфер.

Если вы отписываетесь от одного из GroupedObservables, этот GroupedObservable будет завершён. Если исходный Observable позже излучит элемент, ключ которого соответствует GroupedObservables, который был завершён таким образом, groupBy создаст и излучит новый GroupedObservable для соответствия ключу.

groupBy по умолчанию не работает с каким-либо конкретным Scheduler.

  • Javadoc: groupBy(Func1)
  • Javadoc: groupBy(Func1,Func1)

RxJS groupBy groupByUntil

groupBy

RxJS реализует groupBy. Он принимает от одного до трёх параметров:

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

Пример кода

var codes = [
    { keyCode: 38}, // up
    { keyCode: 38}, // up
    { keyCode: 40}, // down
    { keyCode: 40}, // down
    { keyCode: 37}, // left
    { keyCode: 39}, // right
    { keyCode: 37}, // left
    { keyCode: 39}, // right
    { keyCode: 66}, // b
    { keyCode: 65}  // a
];

var source = Rx.Observable.fromArray(codes)
    .groupBy(
        function (x) { return x.keyCode; },
        function (x) { return x.keyCode; });

var subscription = source.subscribe(
    function (obs) {
        // Print the count
        obs.count().subscribe(function (x) {
            console.log('Count: ' + x);
        });
    },
    function (err) {
        console.log('Error: ' + err);
    },
    function () {
        console.log('Completed');
    });
Count: 2
Count: 2
Count: 2
Count: 2
Count: 1
Count: 1
Completed

groupBy доступен в следующих дистрибутивах:

  • rx.all.js
  • rx.all.compat.js
  • rx.coincidence.js
groupByUntil

RxJS также реализует groupByUntil. Он отслеживает дополнительный Observable, и каждый раз, когда этот Observable излучает элемент, он закрывает все открытые Observable с ключами (он откроет новые, если дополнительные элементы из исходного Observable соответствуют ключу). groupByUntil принимает от двух до четырёх параметров:

  1. (обязательно) функция, которая принимает элемент из исходного Observable и возвращает его ключ
  2. функция, которая принимает элемент из исходного Observable и возвращает элемент, который должен быть излучен вместо него одним из результирующих Observables
  3. (обязательно) функция, возвращающая Observable, излучения которого вызывают завершение любых открытых Observables
  4. функция, используемая для сравнения двух ключей для идентичности (то есть, для того, чтобы определить, должны ли элементы с двумя ключами излучаться в одном Observable)

Пример кода

var codes = [
    { keyCode: 38}, // up
    { keyCode: 38}, // up
    { keyCode: 40}, // down
    { keyCode: 40}, // down
    { keyCode: 37}, // left
    { keyCode: 39}, // right
    { keyCode: 37}, // left
    { keyCode: 39}, // right
    { keyCode: 66}, // b
    { keyCode: 65}  // a
];

var source = Rx.Observable
    .for(codes, function (x) { return Rx.Observable.return(x).delay(1000); })
    .groupByUntil(
        function (x) { return x.keyCode; },
        function (x) { return x.keyCode; },
        function (x) { return Rx.Observable.timer(2000); });

var subscription = source.subscribe(
    function (obs) {
        // Print the count
        obs.count().subscribe(function (x) { console.log('Count: ' + x); });
    },
    function (err) {
        console.log('Error: ' + err);
    },
    function () {
        console.log('Completed');
    });
Count: 2
Count: 2
Count: 1
Count: 1
Count: 1
Count: 1
Count: 1
Count: 1
Completed

groupByUntil доступен в следующих дистрибутивах:

  • rx.all.js
  • rx.all.compat.js
  • rx.coincidence.js

RxPHP groupBy groupByUntil partition

RxPHP реализует этот оператор как groupBy.

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

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/groupBy/groupBy.php

$observable = \Rx\Observable::fromArray([21, 42, 21, 42, 21, 42]);
$observable
    ->groupBy(
        function ($elem) {
            if ($elem === 42) {
                return 0;
            }

            return 1;
        },
        null,
        function ($key) {
            return $key;
        }
    )
    ->subscribe(function ($groupedObserver) use ($createStdoutObserver) {
        $groupedObserver->subscribe($createStdoutObserver($groupedObserver->getKey() . ": "));
    });
1: Next value: 21
0: Next value: 42
1: Next value: 21
0: Next value: 42
1: Next value: 21
0: Next value: 42
1: Complete!
0: Complete!

RxPHP также имеет оператор groupByUntil.

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

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/groupBy/groupByUntil.php

$codes = [
    ['id' => 38],
    ['id' => 38],
    ['id' => 40],
    ['id' => 40],
    ['id' => 37],
    ['id' => 39],
    ['id' => 37],
    ['id' => 39],
    ['id' => 66],
    ['id' => 65]
];

$source = Rx\Observable
    ::fromArray($codes)
    ->concatMap(function ($x) {
        return \Rx\Observable::timer(100)->mapTo($x);
    })
    ->groupByUntil(
        function ($x) {
            return $x['id'];
        },
        function ($x) {
            return $x['id'];
        },
        function ($x) {
            return Rx\Observable::timer(200);
        });

$subscription = $source->subscribe(new CallbackObserver(
    function (\Rx\Observable $obs) {
        // Print the count
        $obs->count()->subscribe(new CallbackObserver(
            function ($x) {
                echo 'Count: ', $x, PHP_EOL;
            }));
    },
    function (Throwable $err) {
        echo 'Error', $err->getMessage(), PHP_EOL;
    },
    function () {
        echo 'Completed', PHP_EOL;
    }));
Count: 2
Count: 2
Count: 1
Count: 1
Count: 1
Count: 1
Count: 1
Count: 1
Completed

RxPHP также имеет оператор partition.

Возвращает два Observable, которые разделяют наблюдения источника по заданной функции. Первый будет активировать наблюдения для тех значений, для которых предикат возвращает true. Второй будет активировать наблюдения для тех значений, для которых предикат возвращает false. Предикат выполняется один раз для каждого подписчика. Оба также распространяют все ошибки, возникающие в источнике, и завершаются, когда источник завершается.

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/partition/partition.php

list($evens, $odds) = \Rx\Observable::range(0, 10, \Rx\Scheduler::getImmediate())
    ->partition(function ($x) {
        return $x % 2 === 0;
    });

//Because we used the immediate scheduler with range, the subscriptions are not asynchronous.
$evens->subscribe($createStdoutObserver('Evens '));
$odds->subscribe($createStdoutObserver('Odds '));
Evens Next value: 0
Evens Next value: 2
Evens Next value: 4
Evens Next value: 6
Evens Next value: 8
Evens Complete!
Odds Next value: 1
Odds Next value: 3
Odds Next value: 5
Odds Next value: 7
Odds Next value: 9
Odds Complete!

© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/groupby.html

Spec-Zone.ru

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