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

JEP 485: Stream Gatherers

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

ОтветственныйViktor Klang
ТипFeature
ОбластьSE
СтатусClosed / Delivered
Выпуск24
Компонентcore-libs / java.util.stream
Обсуждениеcore dash libs dash dev at openjdk dot org
ТрудоёмкостьM
ДлительностьM
Связан сJEP 473: Stream Gatherers (Second Preview)
РецензентыAlan Bateman, Paul Sandoz
ОдобренPaul Sandoz
Создан2024/07/08 15:42
Обновлён2026/01/27 16:01
Задача8335899

Аннотация

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

История

Stream Gatherers были предложены как возможность в статусе Preview (предварительная версия) в JEP 461 в JDK 22 и повторно выпущены в статусе Preview в JEP 473 в JDK 23. Здесь мы предлагаем окончательно утвердить этот API в JDK 24 без изменений.

Цели

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

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

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

  • Изменение языка программирования 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 фиксированного размера, поскольку collector не может сообщить collect, что он завершил работу, пока появляются новые элементы, — а в бесконечном потоке это происходит вечно. Кроме того, задача по своей сути связана с упорядоченными данными, поэтому поручить collector группировку параллельно невозможно, и он должен сообщать об этом, выбрасывая исключение при вызове своего 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 может преобразовывать элементы по схеме «один к одному», «один ко многим», «многие к одному» или «многие ко многим». Он может отслеживать ранее встреченные элементы, чтобы влиять на преобразование последующих, может досрочно завершать работу, чтобы превращать бесконечные потоки в конечные, и может обеспечивать параллельное выполнение. Например, 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(...)

Встроенные gatherers

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

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

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

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

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

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

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

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

Композиция gatherers

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

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

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

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

Gatherers и collectors

Устройство интерфейса Gatherer во многом основано на устройстве Collector. Основные различия:

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

  • Gatherer использует для finisher BiConsumer вместо Function, потому что ему нужен дополнительный входной параметр для объекта Downstream и потому что finisher не может возвращать результат и поэтому имеет тип 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 вычисляется параллельно, только если он предоставляет функцию-комбинатор. Например, этот 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, а также встроенных реализаций, объявленных в классе Gatherers, будет не таким лаконичным, как использование встроенных промежуточных операций, объявленных в классе Stream. Однако определение собственных реализаций gatherer будет сопоставимо по сложности с определением собственных коллекторов для терминальных операций collect. Кроме того, использование как собственных, так и встроенных реализаций gatherer будет сопоставимо по сложности с использованием собственных коллекторов и встроенных коллекторов, объявленных в классе Collectors.

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

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