openjdk.ruOpenJDK на русском

JEP 473: Stream Gatherers (Second Preview)

Stream Gatherers (сборщики промежуточных операций потоков), вторая версия Preview (предварительная версия)

ОтветственныйViktor Klang
ТипFeature
ОбластьSE
СтатусClosed / Delivered
Выпуск23
Компонентcore-libs / java.util.stream
Обсуждениеcore dash libs dash dev at openjdk dot org
ТрудоёмкостьM
ДлительностьM
Связан сJEP 461: Stream Gatherers (Preview)
JEP 485: Stream Gatherers
РецензентыAlan Bateman, Paul Sandoz
ОдобренPaul Sandoz
Создан2024/03/11 19:39
Обновлён2025/06/12 14:49
Задача8327844

Аннотация

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

История

Мы предложили Stream Gatherers как Preview-возможность в JEP 461 и выпустили её в JDK 22. Здесь мы предлагаем повторно выпустить этот API в статусе Preview в JDK 23 без изменений, чтобы получить дополнительный опыт и отзывы.

Цели

  • Сделать конвейеры потоков более гибкими и выразительными.

  • По возможности позволить пользовательским промежуточным операциям работать с потоками бесконечного размера.

Что не является целью

  • Изменение языка программирования Java ради более удобной обработки потоков не является целью.

  • Особая обработка при компиляции кода, использующего Stream API, не является целью.

Мотивация

В Java 8 появился первый API, спроектированный специально для лямбда-выражений: Stream API, java.util.stream. Поток — это лениво вычисляемая, потенциально неограниченная последовательность значений. API позволяет обрабатывать поток как последовательно, так и параллельно.

Конвейер потока состоит из трёх частей: источника элементов, любого числа промежуточных операций и терминальной операции. Например:

long numberOfWords =
    Stream.of("the", "", "fox", "jumps", "over", "the", "", "dog")  // (1)
          .filter(Predicate.not(String::isEmpty))                   // (2)
          .collect(Collectors.counting());                          // (3)

Такой стиль программирования одновременно выразителен и эффективен. В API в стиле builder каждая промежуточная операция возвращает новый поток; вычисление начинается только при вызове терминальной операции. В этом примере строка (1) создаёт поток, но не вычисляет его, строка (2) задаёт промежуточную операцию filter, но по-прежнему не вычисляет поток, и, наконец, терминальная операция collect в строке (3) вычисляет весь конвейер потока.

Stream API предоставляет достаточно богатый, хотя и фиксированный, набор промежуточных и терминальных операций: отображение, фильтрация, свёртка, сортировка и так далее. В нём также есть расширяемая терминальная операция Stream::collect, с помощью которой результат конвейера можно обобщать самыми разными способами.

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

Например, пусть задача — взять поток строк и оставить в нём только различающиеся строки, но различие определять по длине строки, а не по содержимому. То есть должно быть выдано не более одной строки длины 1, не более одной строки длины 2, не более одной строки длины 3 и так далее. В идеале код выглядел бы примерно так:

var result = Stream.of("foo", "bar", "baz", "quux")
                   .distinctBy(String::length)      // Hypothetical
                   .toList();

// result ==> [foo, quux]

К сожалению, distinctBy — не встроенная промежуточная операция. Ближайшая встроенная операция, distinct, отслеживает уже встреченные элементы, сравнивая их с помощью равенства объектов. То есть distinct хранит состояние, но в данном случае не то: нам нужно отслеживать элементы по равенству длины строк, а не их содержимого. Это ограничение можно обойти: объявить класс, который определяет равенство объектов через длину строки, обернуть каждую строку в экземпляр этого класса и применить distinct к этим экземплярам. Однако такое выражение задачи неинтуитивно, а код получается трудным в сопровождении:

record DistinctByLength(String str) {

    @Override public boolean equals(Object obj) {
        return obj instanceof DistinctByLength(String other)
               && str.length() == other.length();
    }

    @Override public int hashCode() {
        return str == null ? 0 : Integer.hashCode(str.length());
    }

}

var result = Stream.of("foo", "bar", "baz", "quux")
                   .map(DistinctByLength::new)
                   .distinct()
                   .map(DistinctByLength::str)
                   .toList();

// result ==> [foo, quux]

Другой пример: пусть задача — сгруппировать элементы в группы фиксированного размера по три, но оставить только первые две группы: [0, 1, 2, 3, 4, 5, 6, ...] должно дать [[0, 1, 2], [3, 4, 5]]. В идеале код выглядел бы так:

var result = Stream.iterate(0, i -> i + 1)
                   .windowFixed(3)                  // Hypothetical
                   .limit(2)
                   .toList();

// result ==> [[0, 1, 2], [3, 4, 5]]

К сожалению, ни одна встроенная промежуточная операция не поддерживает эту задачу. Лучший вариант — поместить логику группировки фиксированными окнами в терминальную операцию, вызвав collect с пользовательским Collector. Однако перед операцией collect приходится поставить операцию limit с фиксированным размером, поскольку коллектор не может сообщить collect, что он завершил работу, пока появляются новые элементы, а в бесконечном потоке это происходит бесконечно. Кроме того, задача по своей сути касается упорядоченных данных, поэтому поручить коллектору группировку в параллельном режиме невозможно, и он должен сообщать об этом, выбрасывая исключение при вызове своей функции combiner. Получившийся код трудно понять:

var result
    = Stream.iterate(0, i -> i + 1)
            .limit(3 * 2)
            .collect(Collector.of(
                () -> new ArrayList<ArrayList<Integer>>(),
                (groups, element) -> {
                    if (groups.isEmpty() || groups.getLast().size() == 3) {
                         var current = new ArrayList<Integer>();
                         current.add(element);
                         groups.addLast(current);
                     } else {
                         groups.getLast().add(element);
                     }
                },
                (left, right) -> {
                    throw new UnsupportedOperationException("Cannot be parallelized");
                }
            ));

// result ==> [[0, 1, 2], [3, 4, 5]]

За прошедшие годы для Stream API было предложено много новых промежуточных операций. Большинство из них имеют смысл по отдельности, но добавление их всех сделало бы (и без того большой) Stream API более сложным в изучении, потому что находить его операции стало бы труднее.

Разработчики Stream API понимали, что желательно иметь точку расширения, чтобы кто угодно мог определять промежуточные операции над потоками. Однако в то время они не знали, как должна выглядеть эта точка расширения. Со временем стало ясно, что точка расширения для терминальных операций, а именно Stream::collect(Collector), оказалась эффективной. Теперь мы можем применить аналогичный подход к промежуточным операциям.

Итак, больше промежуточных операций — больше пользы в конкретных ситуациях, и потоки подходят для ещё большего числа задач. Нам следует предоставить API для пользовательских промежуточных операций, с помощью которого разработчики смогут преобразовывать конечные и бесконечные потоки так, как им удобно.

Описание

Stream::gather(Gatherer) — новая промежуточная операция над потоком, которая обрабатывает элементы потока, применяя определённую пользователем сущность, называемую gatherer. С помощью операции gather можно строить эффективные, готовые к параллельной работе потоки, реализующие почти любую промежуточную операцию. Stream::gather(Gatherer) для промежуточных операций — то же, что Stream::collect(Collector) для терминальных.

Gatherer представляет преобразование элементов потока; это экземпляр интерфейса java.util.stream.Gatherer. Gatherer может преобразовывать элементы по схеме «один к одному», «один ко многим», «многие к одному» или «многие ко многим». Он может отслеживать ранее встреченные элементы, чтобы влиять на преобразование последующих, может завершаться досрочно (short-circuit), чтобы превращать бесконечные потоки в конечные, и может обеспечивать параллельное выполнение. Например, gatherer может преобразовывать один входной элемент в один выходной, пока не станет истинным некоторое условие, после чего начинает преобразовывать один входной элемент в два выходных.

Gatherer определяется четырьмя функциями, работающими совместно:

  • Необязательная функция initializer предоставляет объект, который хранит приватное состояние при обработке элементов потока. Например, gatherer может сохранить текущий элемент, чтобы при следующем применении сравнить новый элемент с уже предыдущим и, скажем, выдать только больший из двух. По сути, такой gatherer преобразует два входных элемента в один выходной.

  • Функция integrator интегрирует новый элемент из входного потока, возможно, проверяя объект приватного состояния и, возможно, выдавая элементы в выходной поток. Она также может завершить обработку до достижения конца входного потока; например, gatherer, ищущий наибольшее из потока целых чисел, может завершить работу, если обнаружит Integer.MAX_VALUE.

  • Необязательную функцию combiner можно использовать для параллельного вычисления gatherer, когда входной поток помечен как параллельный. Если gatherer не поддерживает параллельную работу, он всё равно может быть частью параллельного конвейера потока, но вычисляется последовательно. Это полезно в случаях, когда операция по своей природе упорядочена и поэтому не может быть распараллелена.

  • Необязательная функция finisher вызывается, когда входных элементов для обработки больше нет. Эта функция может проверить объект приватного состояния и, возможно, выдать дополнительные выходные элементы. Например, gatherer, ищущий определённый элемент среди входных, может сообщить о неудаче, скажем, выбросив исключение, при вызове своей функции finisher.

При вызове Stream::gather выполняет действия, эквивалентные следующим шагам:

  • Создать объект Downstream, который, получив элемент выходного типа gatherer, передаёт его на следующую стадию конвейера.

  • Получить объект приватного состояния gatherer, вызвав метод get() его функции initializer.

  • Получить функцию integrator gatherer, вызвав его метод integrator().

  • Пока есть входные элементы, вызывать метод integrate(...) функции integrator, передавая ему объект состояния, следующий элемент и объект downstream. Завершить работу, если этот метод вернёт false.

  • Получить функцию finisher gatherer и вызвать её с объектами состояния и downstream.

Любую существующую промежуточную операцию, объявленную в интерфейсе Stream, можно реализовать, вызвав gather с gatherer, реализующим эту операцию. Например, для потока элементов типа T операция Stream::map превращает каждый элемент T в элемент U, применяя функцию, и затем передаёт элемент U дальше по конвейеру; это просто gatherer без состояния, работающий по схеме «один к одному». Другой пример: Stream::filter принимает предикат, определяющий, нужно ли передавать входной элемент дальше по конвейеру; это просто gatherer без состояния, работающий по схеме «один ко многим». Фактически любой конвейер потока концептуально эквивалентен

source.gather(...).gather(...).gather(...).collect(...)

Встроенные gatherer

Мы добавляем следующие встроенные gatherer в класс java.util.stream.Gatherers:

  • fold — gatherer с состоянием, работающий по схеме «многие к одному», который постепенно строит агрегат и выдаёт его, когда входных элементов больше нет.

  • mapConcurrent — gatherer с состоянием, работающий по схеме «один к одному», который конкурентно вызывает заданную функцию для каждого входного элемента, вплоть до заданного предела.

  • scan — gatherer с состоянием, работающий по схеме «один к одному», который применяет заданную функцию к текущему состоянию и текущему элементу, чтобы получить следующий элемент, и передаёт его дальше по конвейеру.

  • windowFixed — gatherer с состоянием, работающий по схеме «многие ко многим», который группирует входные элементы в списки заданного размера и передаёт окна дальше по конвейеру, когда они заполнены.

  • windowSliding — gatherer с состоянием, работающий по схеме «многие ко многим», который группирует входные элементы в списки заданного размера. После первого окна каждое следующее окно создаётся из копии предыдущего: первый элемент отбрасывается, а в конец добавляется следующий элемент из входного потока..

Параллельное вычисление

Параллельное вычисление gatherer делится на два различных режима. Если combiner не предоставлен, библиотека потоков всё равно может извлечь параллелизм, выполняя операции выше и ниже по конвейеру параллельно, аналогично операции parallel().forEachOrdered() с досрочным завершением. Если combiner предоставлен, параллельное вычисление аналогично операции parallel().reduce() с досрочным завершением.

Композиция gatherer

Gatherer поддерживают композицию с помощью метода andThen(Gatherer), который соединяет два gatherer так, что первый выдаёт элементы, которые может потреблять второй. Так можно создавать сложные gatherer, составляя их из более простых, — точно так же, как при композиции функций. Семантически

source.gather(a).gather(b).gather(c).collect(...)

эквивалентно

source.gather(a.andThen(b).andThen(c)).collect(...)

Gatherer и коллекторы

На устройство интерфейса Gatherer сильно повлияло устройство Collector. Основные различия:

  • Gatherer использует для поэлементной обработки Integrator вместо BiConsumer, потому что ему нужен дополнительный входной параметр для объекта Downstream и потому что он должен возвращать boolean, указывая, следует ли продолжать обработку.

  • Gatherer использует для своей функции finisher BiConsumer вместо Function, потому что ему нужен дополнительный входной параметр для объекта Downstream и потому что он не может вернуть результат и поэтому имеет тип void.

Пример: остаёмся в потоке

Иногда из-за отсутствия подходящей промежуточной операции приходится вычислять поток в список и выполнять логику анализа в цикле. Предположим, например, что у нас есть поток показаний температуры, упорядоченных по времени:

record Reading(Instant obtainedAt, int kelvins) {

    Reading(String time, int kelvins) {
        this(Instant.parse(time), kelvins);
    }

    static Stream<Reading> loadRecentReadings() {
        // In reality these could be read from a file, a database,
        // a service, or otherwise
        return Stream.of(
                new Reading("2023-09-21T10:15:30.00Z", 310),
                new Reading("2023-09-21T10:15:31.00Z", 312),
                new Reading("2023-09-21T10:15:32.00Z", 350),
                new Reading("2023-09-21T10:15:33.00Z", 310)
        );
    }

}

Предположим далее, что мы хотим обнаруживать в этом потоке подозрительные изменения. Это изменения температуры более чем на 30° по Кельвину между двумя последовательными показаниями в пределах пятисекундного окна времени:

boolean isSuspicious(Reading previous, Reading next) {
    return next.obtainedAt().isBefore(previous.obtainedAt().plusSeconds(5))
           && (next.kelvins() > previous.kelvins() + 30
               || next.kelvins() < previous.kelvins() - 30);
}

Для этого нужен последовательный просмотр входного потока, поэтому приходится отказаться от декларативной обработки потоков и реализовать анализ в императивном стиле:

List<List<Reading>> findSuspicious(Stream<Reading> source) {
    var suspicious = new ArrayList<List<Reading>>();
    Reading previous = null;
    boolean hasPrevious = false;
    for (Reading next : source.toList()) {
        if (!hasPrevious) {
            hasPrevious = true;
            previous = next;
        } else {
            if (isSuspicious(previous, next))
                suspicious.add(List.of(previous, next));
            previous = next;
        }
    }
    return suspicious;
}

var result = findSuspicious(Reading.loadRecentReadings());

// result ==> [[Reading[obtainedAt=2023-09-21T10:15:31Z, kelvins=312],
//              Reading[obtainedAt=2023-09-21T10:15:32Z, kelvins=350]],
//             [Reading[obtainedAt=2023-09-21T10:15:32Z, kelvins=350],
//              Reading[obtainedAt=2023-09-21T10:15:33Z, kelvins=310]]]

Однако с помощью gatherer это можно выразить короче:

List<List<Reading>> findSuspicious(Stream<Reading> source) {
    return source.gather(Gatherers.windowSliding(2))
                 .filter(window -> (window.size() == 2
                                    && isSuspicious(window.get(0),
                                                    window.get(1))))
                 .toList();
}

Пример: определение gatherer

Gatherer windowFixed, объявленный в классе Gatherers, можно было бы написать как непосредственную реализацию интерфейса Gatherer:

record WindowFixed<TR>(int windowSize)
    implements Gatherer<TR, ArrayList<TR>, List<TR>>
{

    public WindowFixed {
        // Validate input
        if (windowSize < 1)
            throw new IllegalArgumentException("window size must be positive");
    }

    @Override
    public Supplier<ArrayList<TR>> initializer() {
        // Create an ArrayList to hold the current open window
        return () -> new ArrayList<>(windowSize);
    }

    @Override
    public Integrator<ArrayList<TR>, TR, List<TR>> integrator() {
        // The integrator is invoked for each element consumed
        return Gatherer.Integrator.ofGreedy((window, element, downstream) -> {

            // Add the element to the current open window
            window.add(element);

            // Until we reach our desired window size,
            // return true to signal that more elements are desired
            if (window.size() < windowSize)
                return true;

            // When the window is full, close it by creating a copy
            var result = new ArrayList<TR>(window);

            // Clear the window so the next can be started
            window.clear();

            // Send the closed window downstream
            return downstream.push(result);

        });
    }

    // The combiner is omitted since this operation is intrinsically sequential,
    // and thus cannot be parallelized

    @Override
    public BiConsumer<ArrayList<TR>, Downstream<? super List<TR>>> finisher() {
        // The finisher runs when there are no more elements to pass from
        // the upstream
        return (window, downstream) -> {
            // If the downstream still accepts more elements and the current
            // open window is non-empty, then send a copy of it downstream
            if(!downstream.isRejecting() && !window.isEmpty()) {
                downstream.push(new ArrayList<TR>(window));
                window.clear();
            }
        };
    }
}

Пример использования:

jshell> Stream.of(1,2,3,4,5,6,7,8,9).gather(new WindowFixed(3)).toList()
$1 ==> [[1, 2, 3], [4, 5, 6], [7, 8, 9]]

Пример: gatherer, определённый на месте

Gatherer windowFixed можно было бы также написать на месте, с помощью фабричного метода Gatherer.ofSequential(...):

/**
 * Gathers elements into fixed-size groups. The last group may contain fewer
 * elements.
 * @param windowSize the maximum size of the groups
 * @return a new gatherer which groups elements into fixed-size groups
 * @param <TR> the type of elements the returned gatherer consumes and produces
 */
static <TR> Gatherer<TR, ?, List<TR>> fixedWindow(int windowSize) {

    // Validate input
    if (windowSize < 1)
      throw new IllegalArgumentException("window size must be non-zero");

    // This gatherer is inherently order-dependent,
    // so it should not be parallelized
    return Gatherer.ofSequential(

            // The initializer creates an ArrayList which holds the current
            // open window
            () -> new ArrayList<TR>(windowSize),

            // The integrator is invoked for each element consumed
            Gatherer.Integrator.ofGreedy((window, element, downstream) -> {

                // Add the element to the current open window
                window.add(element);

                // Until we reach our desired window size,
                // return true to signal that more elements are desired
                if (window.size() < windowSize)
                    return true;

                // When window is full, close it by creating a copy
                var result = new ArrayList<TR>(window);

                // Clear the window so the next can be started
                window.clear();

                // Send the closed window downstream
                return downstream.push(result);

            }),

            // The combiner is omitted since this operation is intrinsically sequential,
            // and thus cannot be parallelized

            // The finisher runs when there are no more elements to pass from the upstream
            (window, downstream) -> {
                // If the downstream still accepts more elements and the current
                // open window is non-empty then send a copy of it downstream
                if(!downstream.isRejecting() && !window.isEmpty()) {
                    downstream.push(new ArrayList<TR>(window));
                    window.clear();
                }
            }

    );
}

Пример использования:

jshell> Stream.of(1,2,3,4,5,6,7,8,9).gather(fixedWindow(3)).toList()
$1 ==> [[1, 2, 3], [4, 5, 6], [7, 8, 9]]

Пример: распараллеливаемый gatherer

В параллельном потоке gatherer вычисляется параллельно, только если он предоставляет функцию combiner. Например, этот распараллеливаемый gatherer выдаёт не более одного элемента на основе переданной функции-селектора:

static <TR> Gatherer<TR, ?, TR> selectOne(BinaryOperator<TR> selector) {

    // Validate input
    Objects.requireNonNull(selector, "selector must not be null");

    // Private state to track information across elements
    class State {
        TR value;            // The current best value
        boolean hasValue;    // true when value holds a valid value
    }

    // Use the `of` factory method to construct a gatherer given a set
    // of functions for `initializer`, `integrator`, `combiner`, and `finisher`
    return Gatherer.of(

            // The initializer creates a new State instance
            State::new,

            // The integrator; in this case we use `ofGreedy` to signal
            // that this integerator will never short-circuit
            Gatherer.Integrator.ofGreedy((state, element, downstream) -> {
                if (!state.hasValue) {
                    // The first element, just save it
                    state.value = element;
                    state.hasValue = true;
                } else {
                    // Select which value of the two to save, and save it
                    state.value = selector.apply(state.value, element);
                }
                return true;
            }),

            // The combiner, used during parallel evaluation
            (leftState, rightState) -> {
                if (!leftState.hasValue) {
                    // If no value on the left, return the right
                    return rightState;
                } else if (!rightState.hasValue) {
                    // If no value on the right, return the left
                    return leftState;
                } else {
                    // If both sides have values, select one of them to keep
                    // and store it in the leftState, as that will be returned
                    leftState.value = selector.apply(leftState.value,
                                                     rightState.value);
                    return leftState;
                }
            },

            // The finisher
            (state, downstream) -> {
                // Emit the selected value, if there is one, downstream
                if (state.hasValue)
                    downstream.push(state.value);
            }

    );
}

Пример использования на потоке случайных целых чисел:

jshell> Stream.generate(() -> ThreadLocalRandom.current().nextInt())
              .limit(1000)                   // Take the first 1000 elements
              .gather(selectOne(Math::max))  // Select the largest value seen
              .parallel()                    // Execute in parallel
              .findFirst()                   // Extract the largest value
$1 ==> Optional[99822]

Альтернативы

Мы рассмотрели альтернативы в отдельном проектном документе.

Риски и допущения

  • Использование собственных gatherer и встроенных gatherer, объявленных в классе Gatherers, будет не таким лаконичным, как использование встроенных промежуточных операций, объявленных в классе Stream. Однако определение собственных gatherer будет сопоставимо по сложности с определением собственных коллекторов для терминальных операций collect. Кроме того, использование как собственных, так и встроенных gatherer будет сопоставимо по сложности с использованием собственных коллекторов и встроенных коллекторов, объявленных в классе Collectors.

  • Мы можем пересмотреть набор встроенных gatherer, пока эта возможность находится в статусе Preview, а также можем пересмотреть его в будущих выпусках.

  • Мы не будем добавлять в класс Stream новую промежуточную операцию для каждого из встроенных gatherer, определённых в классе Gatherers, хотя ради единообразия это и соблазнительно. Чтобы класс Stream оставался простым для изучения, мы рассмотрим добавление в него новых промежуточных операций только после того, как опыт покажет, что они полезны в широком круге случаев. Мы можем добавить такие методы в одной из следующих версий Preview или даже после того, как эта возможность станет окончательной. Если открыть доступ к новым встроенным gatherer сейчас, это не помешает позже добавить специальные методы Stream.