JEP 461: Stream Gatherers (Preview)
Stream Gatherers (сборщики промежуточных операций потоков), версия Preview (предварительная версия)
| Ответственный | Viktor Klang |
| Тип | Feature |
| Область | SE |
| Статус | Closed / Delivered |
| Выпуск | 22 |
| Компонент | core-libs / java.util.stream |
| Обсуждение | core dash libs dash dev at openjdk dot org |
| Трудоёмкость | M |
| Длительность | M |
| Связан с | JEP 473: Stream Gatherers (Second Preview) |
| Рецензенты | Alan Bateman, Alex Buckley, Paul Sandoz |
| Одобрен | Paul Sandoz |
| Создан | 2023/10/11 13:08 |
| Обновлён | 2025/06/12 14:48 |
| Задача | 8317955 |
Аннотация
Расширить Stream API поддержкой пользовательских промежуточных операций. Это позволит конвейерам потоков преобразовывать данные способами, которых трудно добиться с помощью существующих встроенных промежуточных операций. Это API в статусе Preview.
Цели
-
Сделать конвейеры потоков более гибкими и выразительными.
-
По мере возможности позволить пользовательским промежуточным операциям работать с потоками бесконечного размера.
Что не является целью
-
Цель не в том, чтобы изменить язык программирования Java ради более удобной обработки потоков.
-
Цель не в том, чтобы особым образом компилировать код, использующий Stream API.
Мотивация
В Java 8 появился первый API, спроектированный специально для лямбда-выражений: Stream API, java.util.stream. Поток — это лениво вычисляемая, потенциально неограниченная последовательность значений. API позволяет обрабатывать поток последовательно или параллельно.
Конвейер потока (stream pipeline) состоит из трёх частей: источника элементов, любого числа промежуточных операций и терминальной операции. Например:
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, ищущий определённый элемент среди входных элементов, может сообщить о неудаче, скажем, выбросив исключение, при вызове своего финализатора.
При вызове Stream::gather выполняет действия, эквивалентные следующим шагам:
-
Создать объект
Downstream, который, получив элемент выходного типа gatherer, передаёт его следующему этапу конвейера. -
Получить объект приватного состояния gatherer, вызвав метод
get()его инициализатора. -
Получить интегратор gatherer, вызвав его метод
integrator(). -
Пока есть входные элементы, вызывать метод интегратора integrate(...), передавая ему объект состояния, следующий элемент и объект downstream. Завершить работу, если этот метод вернёт
false. -
Получить финализатор 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 делится на два различных режима. Если комбинатор не задан, библиотека потоков всё равно может извлечь параллелизм, выполняя предшествующие и последующие операции параллельно, аналогично операции parallel().forEachOrdered() с возможностью досрочного завершения. Если комбинатор задан, параллельное вычисление аналогично операции 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использует для финализатора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, созданный ad hoc
Gatherer windowFixed можно было бы также написать ad hoc с помощью фабричного метода 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]
Альтернативы
Мы рассмотрели альтернативы в отдельном проектном документе.
Риски и допущения
-
Использование пользовательских gatherers и встроенных gatherers, объявленных в классе
Gatherers, будет не таким лаконичным, как использование встроенных промежуточных операций, объявленных в классеStream. Однако определение пользовательских gatherers по сложности будет сопоставимо с определением пользовательских коллекторов для терминальных операцийcollect. Кроме того, использование как пользовательских, так и встроенных gatherers по сложности будет сопоставимо с использованием пользовательских коллекторов и встроенных коллекторов, объявленных в классеCollectors. -
Мы можем пересмотреть набор встроенных gatherers, пока эта возможность находится в статусе Preview, и можем пересмотреть его в будущих выпусках.
-
Мы не будем добавлять в класс
Streamновую промежуточную операцию для каждого из встроенных gatherers, определённых в классеGatherers, хотя ради единообразия это и заманчиво. Чтобы классStreamоставался простым для изучения, мы рассмотрим добавление в него новых промежуточных операций только после того, как опыт покажет, что они широко полезны. Мы можем добавить такие методы в одной из следующих версий Preview или даже после того, как эта возможность станет окончательной. Появление новых встроенных gatherers сейчас не исключает добавления специальных методовStreamпозже.