Пакет java.util.stream
int sum = widgets.stream()
.filter(b -> b.getColor() == RED)
.mapToInt(b -> b.getWeight())
.sum();
Здесь мы используем widgets, a Collection<Widget>, в качестве источника для потока, а затем выполняем фильтрацию, преобразование и сворачивание над потоком для получения суммы весов красных виджетов. (Суммирование — пример операции свертки.)
Ключевая абстракция, введённая в этом пакете, — это поток. Классы Stream, IntStream, LongStream и DoubleStream являются потоками над объектами и примитивными типами int, long и double. Потоки отличаются от коллекций несколькими способами:
- Нет хранения. Поток не является структурой данных, хранящей элементы; вместо этого он передаёт элементы из источника, такого как структура данных, массив, функция-генератор или канал ввода-вывода, через цепочку вычислительных операций.
- Функциональный характер. Операция над потоком производит результат, но не изменяет свой источник. Например, фильтрация
Stream, полученная из коллекции, производит новуюStreamбез отфильтрованных элементов, а не удаляет элементы из исходной коллекции. - Ленивость. Многие операции над потоком, такие как фильтрация, преобразование или удаление дубликатов, могут быть реализованы лениво, открывая возможности для оптимизации. Например, "найдите первую
Stringс тремя последовательными гласными" не нужно проверять все входные строки. Операции над потоками делятся на промежуточные (порождающиеStream) операции и терминальные (порождающие значения или побочные эффекты) операции. Промежуточные операции всегда ленивые. - Возможная неограниченность. Хотя коллекции имеют конечный размер, потоки могут быть неограниченными. Короткие операции, такие как
limit(n)илиfindFirst(), могут позволить вычисления на бесконечных потоках завершиться за конечное время. - Потребляемость. Элементы потока посещаются только один раз в течение жизни потока. Как и
Iterator, для повторного посещения тех же элементов источника необходимо сгенерировать новый поток.
- Из
Collectionчерез методыstream()иparallelStream(); - Из массива через
Arrays.stream(Object[]); - Из статических фабричных методов классов потоков, таких как
Stream.of(Object[]),IntStream.range(int, int)илиStream.iterate(Object, UnaryOperator); - Строки файла можно получить из
BufferedReader.lines(); - Потоки путей файлов могут быть получены из методов в
Files; - Потоки случайных чисел могут быть получены из
Random.ints(); - Множество других методов, порождающих потоки, в JDK, включая
BitSet.stream(),Pattern.splitAsStream(java.lang.CharSequence)иJarFile.stream().
Операции и цепочки потоков
Операции над потоком делятся на промежуточные и терминальные операции и объединяются в цепочки потоков. Цепочка потоков состоит из источника (например, Collection, массива, функции-генератора или канала ввода-вывода); за которым следуют ноль или более промежуточных операций, таких как Stream.filter или Stream.map; и терминальной операции, такой как Stream.forEach или Stream.reduce.
Промежуточные операции возвращают новый поток. Они всегда ленивые; выполнение промежуточной операции, такой как filter(), на самом деле не выполняет фильтрацию, а вместо этого создаёт новый поток, который при обходе содержит элементы исходного потока, соответствующие заданному предикату. Обход источника цепочки не начинается до выполнения терминальной операции цепочки.
Терминальные операции, такие как Stream.forEach или IntStream.sum, могут обходить поток для получения результата или побочного эффекта. После выполнения терминальной операции цепочка потоков считается потреблённой и больше не может быть использована; если вам нужно повторно обходить тот же источник данных, нужно вернуться к источнику данных, чтобы получить новый поток. Почти во всех случаях терминальные операции являются жадными, завершая обход источника данных и обработку цепочки перед возвратом. Только терминальные операции iterator() и spliterator() не являются; они предоставляются как «безопасный выход», чтобы позволить произвольному управлению клиентом обходом цепочки в случае, если существующие операции не подходят для задачи.
Обработка потоков лениво позволяет значительно повысить эффективность; в цепочке, такой как пример фильтрации-преобразования-суммирования выше, фильтрация, преобразование и суммирование могут быть объединены в один проход по данным с минимальным промежуточным состоянием. Ленивость также позволяет избежать проверки всех данных, когда это не нужно; для операций, таких как "найдите первую строку, длиннее 1000 символов", достаточно проверить только необходимое количество строк, чтобы найти одну, которая имеет нужные характеристики, не проверяя все доступные строки из источника. (Это поведение становится ещё важнее, когда входной поток является бесконечным, а не просто большим.)
Промежуточные операции далее делятся на бессостоятельные и состоятельные операции. Бессостоятельные операции, такие как filter и map, не сохраняют состояние от ранее увиденных элементов при обработке нового элемента — каждый элемент может быть обработан независимо от операций над другими элементами. Состоятельные операции, такие как distinct и sorted, могут включать состояние от ранее увиденных элементов при обработке новых элементов.
Состоятельные операции могут потребовать обработки всех входных данных перед получением результата. Например, нельзя получить никакие результаты от сортировки потока до тех пор, пока не будут видны все элементы потока. В результате при параллельном вычислении некоторые цепочки, содержащие состоятельные промежуточные операции, могут потребовать несколько проходов по данным или потребуют буферизации значительных данных. Цепочки, содержащие только бессостоятельные промежуточные операции, могут обрабатываться за один проход, независимо от того, последовательный он или параллельный, с минимальным буферированием данных.
Кроме того, некоторые операции считаются короткозамыкающими операциями. Промежуточная операция является короткозамыкающей, если при представлении бесконечного ввода она может произвести конечный поток в результате. Терминальная операция является короткозамыкающей, если при представлении бесконечного ввода она может завершиться за конечное время. Наличие короткозамыкающей операции в цепочке является необходимым, но не достаточным условием для нормального завершения обработки бесконечного потока за конечное время.
Параллельность
Обработка элементов с явным циклом for- по своей природе последовательна. Потоки облегчают параллельное выполнение, переформулировав вычисление как цепочку агрегатных операций, а не как императивные операции над каждым отдельным элементом. Все операции над потоками могут выполняться последовательно или параллельно. Реализации потоков в JDK создают последовательные потоки, если параллельность не запрошена явно. Например, Collection имеет методы Collection.stream() и Collection.parallelStream(), которые производят последовательные и параллельные потоки соответственно; другие методы, порождающие потоки, такие как IntStream.range(int, int), производят последовательные потоки, но эти потоки могут быть эффективно распараллелены путём вызова их метода BaseStream.parallel(). Для выполнения предыдущего запроса «суммы весов виджетов» параллельно нужно сделать следующее:
int sumOfWeights = widgets.parallelStream()
.filter(b -> b.getColor() == RED)
.mapToInt(b -> b.getWeight())
.sum();
Единственное различие между последовательными и параллельными версиями этого примера заключается в создании начального потока, используя «parallelStream()» вместо «stream()». Цепочка потоков выполняется последовательно или параллельно в зависимости от режима потока, на котором вызывается терминальная операция. Последовательный или параллельный режим потока можно определить с помощью метода BaseStream.isParallel(), а режим потока можно изменить с помощью операций BaseStream.sequential() и BaseStream.parallel(). Последнее установленное значение последовательного или параллельного режима применяется к выполнению всей цепочки потоков.
За исключением операций, явно обозначенных как недетерминированные, таких как findAny(), последовательное или параллельное выполнение потока не должно изменять результат вычисления.
Большинство операций над потоками принимают параметры, описывающие поведение пользователя, которые часто являются лямбда-выражениями. Для сохранения правильного поведения эти параметры поведения должны быть не вмешивающимися, а в большинстве случаев — бессостоятельными. Такие параметры всегда являются экземплярами функционального интерфейса, такого как Function, и часто являются лямбда-выражениями или ссылками на методы.
Невмешательство
Потоки позволяют выполнять, возможно, параллельные агрегированные операции над различными источниками данных, включая даже не потокобезопасные коллекции, такие какArrayList. Это возможно только в том случае, если мы можем предотвратить вмешательство в источник данных во время выполнения цепочки потоков. За исключением операций «безопасного выхода» iterator() и spliterator(), выполнение начинается при вызове терминальной операции и заканчивается при завершении терминальной операции. Для большинства источников данных предотвращение вмешательства означает обеспечение того, что источник данных вовсе не изменяется во время выполнения цепочки потоков. Заметное исключение — это потоки, чьи источники являются конкурирующими коллекциями, которые специально разработаны для обработки одновременного изменения. Конкурентные источники потоков — это те, чьи Spliterator сообщают о характеристике CONCURRENT. Соответственно, поведенческие параметры в потоковых конвейерах, источник которых может не быть конкурентным, никогда не должны изменять исходный источник данных потока. Говорят, что поведенческий параметр вмешивается в неконкурентный источник данных, если он изменяет или вызывает изменение источника данных потока. Необходимость невмешательства относится ко всем конвейерам, а не только к параллельным. Если источник потока не является конкурентным, изменение источника данных потока во время выполнения потокового конвейера может привести к исключениям, неверным ответам или несоответствующему поведению. Для хорошо себя ведущих источников потоков источник может быть изменен до начала терминальной операции, и эти изменения будут отражены в покрытых элементах. Например, рассмотрим следующий код:
List<String> l = new ArrayList(Arrays.asList("one", "two"));
Stream<String> sl = l.stream();
l.add("three");
String s = sl.collect(joining(" "));
Сначала создается список, состоящий из двух строк: "one" и "two". Затем из этого списка создается поток. Далее список изменяется путем добавления третьей строки: "three". Наконец, элементы потока собираются и объединяются в строку. Поскольку список был изменен до начала терминальной collect операции, результатом будет строка "one two three". Все потоки, возвращаемые из коллекций JDK, и большинство других классов JDK, ведут себя таким образом; для потоков, сгенерированных другими библиотеками, см. Создание потоков низкого уровня для требований к созданию хорошо себя ведущих потоков. Бессостояние поведения
Результаты потокового конвейера могут быть недетерминированными или неверными, если поведенческие параметры потоковых операций являются состоятельными. Состоятельная лямбда-функция (или другой объект, реализующий соответствующий функциональный интерфейс) — это функция, результат которой зависит от состояния, которое может измениться во время выполнения потокового конвейера. Пример состоятельной лямбда-функции — параметр дляmap() в:
Set<Integer> seen = Collections.synchronizedSet(new HashSet<>());
stream.parallel().map(e -> { if (seen.add(e)) return 0; else return e; })...
В данном случае, если операция отображения выполняется параллельно, результаты для одного и того же входного значения могут отличаться от выполнения к выполнению из-за различий в планировании потоков, тогда как со статической лямбда-выражением результаты всегда будут одинаковыми. Обратите также внимание, что попытка доступа к изменяемому состоянию из поведенческих параметров представляет собой плохой выбор с точки зрения безопасности и производительности; если вы не синхронизируете доступ к этому состоянию, у вас возникает гонка данных, и поэтому ваш код не работает, но если вы синхронизируете доступ к этому состоянию, вы рискуете столкнуться с конкуренцией, которая подрывает параллелизм, от которого вы хотите извлечь выгоду. Лучший подход — полностью избегать состоятельных поведенческих параметров для потоковых операций; обычно есть способ перестроить потоковый конвейер, чтобы избежать состоятельности.
Побочные эффекты
Побочные эффекты в поведенческих параметрах потоковых операций, как правило, не рекомендуется, так как они часто могут привести к непреднамеренным нарушениям требования бессостоятельности, а также к другим проблемам с безопасностью потоков.Если поведенческие параметры имеют побочные эффекты, за исключением явно указанного, нет никаких гарантий относительно:
- видимости этих побочных эффектов для других потоков видимость;
- того, что разные операции над одним и тем же элементом в одном и том же потоковом конвейере выполняются в одном и том же потоке; и
- того, что поведенческие параметры всегда вызываются, так как реализация потока может исключить операции (или целые этапы) из потокового конвейера, если она может доказать, что это не повлияет на результат вычисления.
Порядок побочных эффектов может быть неожиданным. Даже когда конвейер ограничен получением результата, который согласуется с порядком встречи источника потока (например, IntStream.range(0,5).parallel().map(x -> x*2).toArray() должен генерировать [0, 2, 4, 6, 8]), не гарантируется порядок, в котором функция отображения применяется к отдельным элементам, или в каком потоке любой поведенческий параметр выполняется для данного элемента.
Исключение побочных эффектов также может быть неожиданным. За исключением терминальных операций forEach и forEachOrdered, побочные эффекты поведенческих параметров могут не всегда выполняться, когда реализация потока может оптимизировать выполнение поведенческих параметров, не затрагивая результат вычисления. (Для конкретного примера см. примечание к API, документированное в count операции.)
Многие вычисления, где можно использовать побочные эффекты, могут быть выражены более безопасно и эффективно без побочных эффектов, например, используя редукцию вместо изменяемых аккумуляторов. Однако побочные эффекты, такие как использование println() для отладки, обычно безвредны. Небольшое количество потоковых операций, таких как forEach() и peek(), могут работать только через побочные эффекты; их следует использовать с осторожностью.
В качестве примера того, как преобразовать потоковый конвейер, который неудачно использует побочные эффекты, в конвейер, который их не использует, следующий код ищет в потоке строк те, которые соответствуют заданному регулярному выражению, и помещает совпадения в список.
ArrayList<String> results = new ArrayList<>();
stream.filter(s -> pattern.matcher(s).matches())
.forEach(s -> results.add(s)); // Unnecessary use of side-effects!
Этот код необоснованно использует побочные эффекты. Если его выполнить параллельно, отсутствие потоковой безопасности в ArrayList приведет к неверным результатам, а добавление необходимой синхронизации приведет к конкуренции, что подрывает преимущество параллелизма. Кроме того, использование побочных эффектов здесь совершенно не нужно; forEach() можно просто заменить операцией редукции, которая более безопасна, эффективна и более подходит для параллелизации:
List<String> results =
stream.filter(s -> pattern.matcher(s).matches())
.toList(); // No side-effects!
Порядок
Потоки могут или не могут иметь определенный порядок встречи. Наличие порядка встречи потока зависит от источника и промежуточных операций. Некоторые источники потоков (например, List или массивы) изначально упорядочены, тогда как другие (например, HashSet) нет. Некоторые промежуточные операции, такие как sorted(), могут накладывать порядок встречи на поток, который не был упорядочен, а другие могут сделать упорядоченный поток неупорядоченным, например, BaseStream.unordered(). Кроме того, некоторые терминальные операции могут игнорировать порядок встречи, например, forEach().
Если поток упорядочен, большинство операций ограничены работой с элементами в порядке их встречи; если источник потока является List, содержащим [1, 2, 3], то результат выполнения map(x -> x*2) должен быть [2, 4, 6]. Однако, если у источника нет определенного порядка встречи, любая перестановка значений [2, 4, 6] будет допустимым результатом.
Для последовательных потоков наличие или отсутствие порядка встречи не влияет на производительность, а только на детерминизм. Если поток упорядочен, повторное выполнение идентичных потоковых конвейеров на идентичном источнике приведет к идентичному результату; если он не упорядочен, повторное выполнение может привести к разным результатам.
Для параллельных потоков ослабление ограничения порядка иногда может повысить эффективность выполнения. Некоторые агрегатные операции, такие как фильтрация дубликатов (distinct()) или групповые редукции (Collectors.groupingBy()), могут быть реализованы более эффективно, если порядок элементов не имеет значения. Аналогично, операции, которые по своей природе связаны с порядком встречи, например, limit(), могут потребовать буферизации для обеспечения правильного порядка, что подрывает выгоду от параллелизма. В тех случаях, когда у потока есть порядок встречи, но пользователь не особенно заинтересован в этом порядке встречи, явное разбиение порядка потока с помощью unordered() может улучшить параллельную производительность для некоторых состоятельных или терминальных операций. Однако большинство потоковых конвейеров, таких как пример "суммы веса блоков" выше, по-прежнему эффективно паралелизуются даже при ограничениях порядка.
Операции редукции
Операция редукции (также называемая сведением) принимает последовательность входных элементов и объединяет их в один итоговый результат путем многократного применения операции объединения, такой как нахождение суммы или максимума набора чисел, или накопление элементов в список. Классы потоков имеют несколько форм общих операций редукции, называемыхreduce() и collect(), а также несколько специализированных форм редукции, таких как sum(), max() или count(). Конечно, такие операции легко реализуются с помощью простых последовательных циклов, как в:
int sum = 0;
for (int x : numbers) {
sum += x;
}
Однако есть веские причины предпочесть операцию редукции изменяемому накоплению, подобному приведенному выше. Операция редукции не только «более абстрактна» — она работает со всем потоком, а не с отдельными элементами — но и правильно построенная операция редукции по своей природе может быть паралелизована, если используемые функции обработки элементов являются ассоциативными и бессостоятельными. Например, для потока чисел, для которого мы хотим найти сумму, мы можем написать:
int sum = numbers.stream().reduce(0, (x,y) -> x+y);
или:
int sum = numbers.stream().reduce(0, Integer::sum);
Эти операции редукции могут безопасно выполняться параллельно практически без модификаций:
int sum = numbers.parallelStream().reduce(0, Integer::sum);
Редукция хорошо паралелизуется, потому что реализация может работать с подмножествами данных параллельно, а затем объединять промежуточные результаты для получения окончательного правильного ответа. (Даже если язык имел конструкцию «параллельный foreach», подход изменяемого накопления все равно потребовал бы от разработчика обеспечить безопасность потоков при обновлении общего переменной накопления sum, а необходимая синхронизация, скорее всего, устранила бы любой выигрыш в производительности от параллелизма.) Использование reduce() вместо этого снимает всю нагрузку по параллелизации операции редукции, и библиотека может предоставить эффективную параллельную реализацию без дополнительной синхронизации.
Примеры "виджетов", показанные ранее, демонстрируют, как редукция сочетается с другими операциями для замены циклов for операциями по обработке массивов. Если widgets — это коллекция объектов Widget, у которых есть метод getWeight, мы можем найти самый тяжелый виджет с помощью:
OptionalInt heaviest = widgets.parallelStream()
.mapToInt(Widget::getWeight)
.max();
В своей наиболее общей форме операция reduce над элементами типа <T>, дающая результат типа <U>, требует трех параметров:
<U> U reduce(U identity,
BiFunction<U, ? super T, U> accumulator,
BinaryOperator<U> combiner);
Здесь единичный элемент является и начальным значением для редукции, и значением по умолчанию, если нет входных элементов. Функция аккумулятора принимает частичный результат и следующий элемент и генерирует новый частичный результат. Функция комбинированияБолее формально, значение identity должно быть тождеством для функции-комбинатора. Это означает, что для всех u, combiner.apply(identity, u) равно u. Кроме того, функция combiner должна быть ассоциативной и должна быть совместима с функцией accumulator: для всех u и t, combiner.apply(u, accumulator.apply(identity, t)) должна быть equals() с accumulator.apply(u, t).
Трех-аргументная форма обобщает двух-аргументную форму, включающую этап отображения в этап накопления. Мы могли бы переформулировать пример суммы весов, используя более общую форму следующим образом:
int sumOfWeights = widgets.stream()
.reduce(0,
(sum, b) -> sum + b.getWeight(),
Integer::sum);
хотя явная форма map-reduce более читаема и поэтому обычно предпочтительнее. Обобщённая форма предоставляется для случаев, когда значительная работа может быть оптимизирована путём объединения отображения и сокращения в одну функцию. Мутабельное сокращение
Мутабельная операция сокращения накапливает элементы входных данных в мутабельный контейнер результата, например, вCollection или StringBuilder, по мере обработки элементов в потоке. Если мы хотели бы взять поток строк и объединить их в одну длинную строку, мы могли бы достичь этого с помощью обычного сокращения:
String concatenated = strings.reduce("", String::concat)
Мы бы получили желаемый результат, и это даже работало бы в параллельном режиме. Однако нас, возможно, не устроит производительность! Такая реализация выполнила бы большое количество копирования строк, а время выполнения было бы O(n^2) в количестве символов. Более производительный подход состоял бы в накоплении результатов в StringBuilder, который является мутабельным контейнером для накопления строк. Мы можем использовать ту же технику для распараллеливания мутабельного сокращения, что и для обычного сокращения.
Операция мутабельного сокращения называется collect(), поскольку она собирает желаемые результаты в контейнер результата, такой как Collection. Операция collect требует трех функций: функцию поставщика для создания новых экземпляров контейнера результата, функцию накопителя для включения элемента входных данных в контейнер результата и функцию объединения для слияния содержимого одного контейнера результата в другой. Форма этого очень похожа на общую форму обычного сокращения:
<R> R collect(Supplier<R> supplier,
BiConsumer<R, ? super T> accumulator,
BiConsumer<R, R> combiner);
Как и в случае с reduce(), преимущество выражения collect таким абстрактным способом заключается в том, что оно напрямую подходит для распараллеливания: мы можем накапливать частичные результаты параллельно, а затем объединять их, если функции накопления и объединения удовлетворяют соответствующим требованиям. Например, для сбора строковых представлений элементов в потоке в ArrayList мы могли бы написать очевидную последовательную форму for-each:
ArrayList<String> strings = new ArrayList<>();
for (T element : stream) {
strings.add(element.toString());
}
Или мы могли бы использовать форму собирания с возможностью распараллеливания:
ArrayList<String> strings = stream.collect(() -> new ArrayList<>(),
(c, e) -> c.add(e.toString()),
(c1, c2) -> c1.addAll(c2));
или, вынося операцию отображения из функции накопителя, мы могли бы выразить её более кратко, как:
List<String> strings = stream.map(Object::toString)
.collect(ArrayList::new, ArrayList::add, ArrayList::addAll);
Здесь наш поставщик - просто ArrayList constructor, накопитель добавляет строковое представление элемента в ArrayList, а комбинатор просто использует addAll для копирования строк из одного контейнера в другой. Три аспекта collect — поставщик, накопитель и комбинатор — тесно связаны. Мы можем использовать абстракцию Collector, чтобы захватить все три аспекта. Вышеупомянутый пример сбора строк в List может быть переписан с использованием стандартного Collector, как:
List<String> strings = stream.map(Object::toString)
.collect(Collectors.toList());
Упаковка мутабельных сокращений в Collector имеет еще одно преимущество: композиционность. Класс Collectors содержит ряд предопределенных фабрик для коллекторов, включая комбинаторы, которые преобразуют один коллектор в другой. Например, предположим, что у нас есть коллектор, который вычисляет сумму заработных плат потока сотрудников, как показано ниже:
Collector<Employee, ?, Integer> summingSalaries
= Collectors.summingInt(Employee::getSalary);
(? для второго типа параметра просто указывает, что нас не интересует промежутовое представление, используемое этим коллектором.) Если мы хотели создать коллектор для табуляции суммы заработной платы по отделам, мы могли бы повторно использовать summingSalaries, используя groupingBy:
Map<Department, Integer> salariesByDept
= employees.stream().collect(Collectors.groupingBy(Employee::getDepartment,
summingSalaries));
Как и в случае с обычной операцией сокращения, операции collect() могут быть распараллелены только в том случае, если выполнены соответствующие условия. Для любого частично накопленного результата, объединение его с пустым контейнером результата должно давать эквивалентный результат. То есть, для частично накопленного результата p, который является результатом любой последовательности вызовов накопителя и комбинатора, p должен быть эквивалентен combiner.apply(p, supplier.get()).
Кроме того, независимо от того, как вычисление разделено, оно должно давать эквивалентный результат. Для любых входных элементов t1 и t2, результаты r1 и r2 в вычислении ниже должны быть эквивалентны:
A a1 = supplier.get();
accumulator.accept(a1, t1);
accumulator.accept(a1, t2);
R r1 = finisher.apply(a1); // result without splitting
A a2 = supplier.get();
accumulator.accept(a2, t1);
A a3 = supplier.get();
accumulator.accept(a3, t2);
R r2 = finisher.apply(combiner.apply(a2, a3)); // result with splitting
Здесь эквивалентность обычно означает в соответствии с Object.equals(Object), но в некоторых случаях эквивалентность может быть ослаблена, чтобы учесть различия в порядке.
Расширяемость
Реализация Collector; использование метода фабрики java.util.stream.Collector.of(...); или использование предопределенных коллекторов в Collectors позволяет создавать пользовательские, повторно используемые конечные операции.
Реализация Gatherer; использование методов фабрик java.util.stream.Gatherer.of(...) и java.util.stream.Gatherer.ofSequential(...); или использование предопределенных собирателей в Gatherers позволяет создавать пользовательские, повторно используемые промежуточные операции.
Сокращение, одновременность и порядок
При некоторых сложных операциях сокращения, например,collect(), которая производит Map, например:
Map<Buyer, List<Transaction>> salesByBuyer
= txns.parallelStream()
.collect(Collectors.groupingBy(Transaction::getBuyer));
, параллельное выполнение операции может быть неэффективным. Это происходит потому, что этап объединения (слияние одного Map в другой по ключу) может быть дорогим для некоторых реализаций Map. Однако, если контейнер результата, используемый в этом сокращении, является одновременно изменяемой коллекцией — например, ConcurrentHashMap, — параллельные вызовы накопителя фактически могут помещать свои результаты одновременно в тот же общий контейнер результата, устраняя необходимость для комбинатора в слиянии различных контейнеров результата. Это потенциально повышает производительность параллельного выполнения. Мы называем это одновременным сокращением.
Collector Collector, поддерживающий одновременное сокращение, помечается характеристикой Collector.Characteristics.CONCURRENT. Однако одновременная коллекция также имеет недостаток. Если несколько потоков одновременно помещают результаты в общий контейнер, порядок, в котором результаты размещаются, не определён. Следовательно, одновременное сокращение возможно только в том случае, если порядок не важен для обрабатываемого потока. Реализация Stream.collect(Collector) выполнит одновременное сокращение только если
- поток параллелен;
- коллектор имеет характеристику
Collector.Characteristics.CONCURRENT, и; - либо поток неупорядочен, либо коллектор имеет характеристику
Collector.Characteristics.UNORDERED.
BaseStream.unordered(). Например:
Map<Buyer, List<Transaction>> salesByBuyer
= txns.parallelStream()
.unordered()
.collect(groupingByConcurrent(Transaction::getBuyer));
(где Collectors.groupingByConcurrent(java.util.function.Function<? super T, ? extends K>) — это одновременный эквивалент groupingBy). Обратите внимание, что если важно, чтобы элементы для данного ключа появлялись в том порядке, в котором они появляются в источнике, то мы не можем использовать одновременное сокращение, поскольку порядок является одной из жертв одновременной вставки. Тогда мы были бы ограничены реализацией либо последовательного сокращения, либо основанного на слиянии параллельного сокращения.
Ассоциативность
Оператор или функцияop является ассоциативной, если выполняется следующее:
(a op b) op c == a op (b op c)
Важность этого для параллельной оценки можно увидеть, если расширить это до четырёх терминов:
a op b op c op d == (a op b) op (c op d)
Таким образом, мы можем вычислить (a op b) параллельно с (c op d), а затем вызвать op на результатах. Примеры ассоциативных операций включают числовое сложение, min, max и конкатенацию строк.
Создание потоков низкого уровня
До сих пор все примеры потоков использовали методы, такие какCollection.stream() или Arrays.stream(Object[]) для получения потока. Как реализуются эти методы, возвращающие поток? Класс StreamSupport имеет ряд методов низкого уровня для создания потока, все они используют какую-то форму Spliterator. Spliterator — это аналог Iterator для параллельного случая; он описывает (возможно, бесконечную) коллекцию элементов с поддержкой последовательного продвижения, объёмной обработки и разделения части входных данных на другой spliterator, который может обрабатываться параллельно. На низшем уровне все потоки управляются spliterator.
Существует ряд вариантов реализации разделителя, почти все из которых представляют собой компромисс между простотой реализации и производительностью потоков, использующих этот разделитель. Самый простой, но наименее производительный способ создания разделителя — создание его из итератора с помощью Spliterators.spliteratorUnknownSize(java.util.Iterator, int). Хотя такой разделитель будет работать, он, вероятно, будет плохо выполнять параллельную обработку, так как мы потеряли информацию о размере (каков размер исходного набора данных), а также ограничены простым алгоритмом разделения.
Разделитель более высокого качества обеспечит сбалансированные и известные разделы, точную информацию о размере и ряд других characteristics разделителя или данных, которые могут использоваться реализациями для оптимизации выполнения.
Разделители для изменяемых источников данных имеют дополнительную проблему: время привязки к данным, так как данные могут измениться между временем создания разделителя и временем выполнения потоковой обработки. В идеале, разделитель для потока должен сообщать о характеристике IMMUTABLE или CONCURRENT; в противном случае он должен быть поздне-связанным. Если источник не может напрямую предоставить рекомендуемый разделитель, он может косвенно предоставить разделитель, используя Supplier, и создать поток через версии stream(), принимающие Supplier. Разделитель извлекается из поставщика только после запуска завершающей операции потоковой обработки.
Эти требования значительно сокращают объем потенциального конфликта между изменениями источника потока и выполнением потоковых конвейеров. Потоки, основанные на разделителях с желаемыми характеристиками или использующих поставщик-ориентированные формы фабрики, невосприимчивы к изменениям источника данных до начала завершающей операции (при условии, что параметры поведения операций потока соответствуют требуемым критериям отсутствия конфликта и бессостоятельности). Подробнее см. Непосредственное взаимодействие.
- С:
- 1.8
| Класс | Описание |
|---|---|
|
BaseStream<T, S extends BaseStream<T, |
Основной интерфейс для потоков, которые представляют собой последовательности элементов, поддерживающие последовательные и параллельные агрегированные операции. |
|
Collector<T, |
Операция изменяемого сокращения, которая накапливает входные элементы в изменяемом контейнере результатов, необязательно преобразуя накопленный результат в окончательное представление после обработки всех входных элементов. |
| Collector.Characteristics | Характеристики, указывающие свойства Collector, которые могут использоваться для оптимизации реализации сокращения. |
| Collectors | Реализации Collector, которые реализуют различные полезные операции сокращения, такие как накопление элементов в коллекции, обобщение элементов по различным критериям и т.д. |
| DoubleStream | Последовательность примитивных элементов типа double, поддерживающих последовательные и параллельные агрегированные операции. |
| DoubleStream.Builder | Изменяемый билдер для DoubleStream. |
| DoubleStream.DoubleMapMultiConsumer | Представляет операцию, которая принимает аргумент типа double и DoubleConsumer, и не возвращает результат. |
|
Gatherer<T, |
Промежуточная операция, преобразующая поток входных элементов в поток выходных элементов, необязательно применяя конечное действие при достижении конца входного потока. |
| Gatherer.Downstream<T> | Объект Downstream — это следующая стадия в конвейере операций, куда можно отправлять элементы. |
|
Gatherer.Integrator<A, |
Интегратор получает элементы и обрабатывает их, необязательно используя предоставленное состояние и необязательно отправляет промежуточные результаты в нисходящий поток. |
|
Gatherer.Integrator.Greedy<A, |
Жадные интеграторы потребляют все свои входные данные и могут только сообщать, что нисходящий поток не хочет больше элементов. |
| Gatherers | Реализации Gatherer, которые предоставляют полезные промежуточные операции, такие как функции окон, функции сворачивания, преобразование элементов параллельно и т.д. |
| IntStream | Последовательность примитивных элементов типа int, поддерживающих последовательные и параллельные агрегированные операции. |
| IntStream.Builder | Изменяемый билдер для IntStream. |
| IntStream.IntMapMultiConsumer | Представляет операцию, которая принимает аргумент типа int и IntConsumer, и не возвращает результат. |
| LongStream | Последовательность примитивных элементов типа long, поддерживающих последовательные и параллельные агрегированные операции. |
| LongStream.Builder | Изменяемый билдер для LongStream. |
| LongStream.LongMapMultiConsumer | Представляет операцию, которая принимает аргумент типа long и LongConsumer, и не возвращает результат. |
| Stream<T> | Последовательность элементов, поддерживающих последовательные и параллельные агрегированные операции. |
| Stream.Builder<T> | Изменяемый билдер для Stream. |
| StreamSupport | Вспомогательные методы низкого уровня для создания и управления потоками. |
© 1993, 2025, Oracle and/or its affiliates. All rights reserved.
Documentation extracted from Debian's OpenJDK Development Kit package.
Licensed under the GNU General Public License, version 2, with the Classpath Exception.
Various third party code in OpenJDK is licensed under different licenses (see Debian package).
Java and OpenJDK are trademarks or registered trademarks of Oracle and/or its affiliates.
https://download.java.net/java/early_access/jdk24/docs/api/java.base/java/util/stream/package-summary.html