Пакет java.util.stream
int sum = widgets.stream()
.filter(b -> b.getColor() == RED)
.mapToInt(b -> b.getWeight())
.sum();
Здесь мы используем widgets, а 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);
Свертка параллелизуется хорошо, потому что реализация может обрабатывать подмножества данных параллельно, а затем объединять промежуточные результаты, чтобы получить окончательный правильный ответ. (Даже если язык имел бы конструкцию "параллельный цикл for-each", подход с накоплением с изменением всё равно потребовал бы от разработчика предоставления безопасных для потоков обновлений для общей переменной накопления 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);
хотя явное представление «отображение-сведение» более удобочитаемо и поэтому обычно предпочтительнее. Обобщенная форма предоставляется для случаев, когда значительная работа может быть оптимизирована за счет объединения отображения и сведения в одну функцию. Изменяемое сведение
Операция изменяемого сведения накапливает элементы входных данных в изменяемый контейнер результатов, такой как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, мы могли бы написать очевидную последовательную форму «для каждого»:
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), но в некоторых случаях эквивалентность может быть ослаблена для учета различий в порядке.
Сведение, конкурентность и порядок
При некоторых сложных операциях сведения, например,collect(), который производит Map, такой как:
Map<Buyer, List<Transaction>> salesByBuyer
= txns.parallelStream()
.collect(Collectors.groupingBy(Transaction::getBuyer));
на самом деле может быть неэффективно выполнять операцию параллельно. Это связано с тем, что этап комбинирования (объединение одного Map с другим по ключу) может быть дорогостоящим для некоторых реализаций Map. Однако, если контейнер результатов, используемый в этом сведении, был изменяемым по ключу — например, ConcurrentHashMap. В этом случае параллельные вызовы накопителя могут фактически параллельно помещать свои результаты в тот же общий контейнер результатов, устраняя необходимость для комбинирования объединять разные контейнеры результатов. Это потенциально повышает производительность параллельного выполнения. Мы называем это конкурентным сведением.
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.
Существует ряд вариантов реализации spliterator, почти все из которых представляют собой компромисс между простотой реализации и временем выполнения потоков, использующих этот spliterator. Самый простой, но наименее производительный способ создания spliterator — создание его из итератора с помощью Spliterators.spliteratorUnknownSize(java.util.Iterator, int). Хотя такой spliterator будет работать, он, скорее всего, будет иметь низкую производительность при параллельном использовании, поскольку мы потеряли информацию о размере (каков размер базового набора данных), а также ограничены простым алгоритмом разделения.
Spliterator более высокого качества обеспечит сбалансированные и известные разделы, точную информацию о размере и ряд других characteristics spliterator или данных, которые могут быть использованы реализациями для оптимизации выполнения.
У spliterators для изменяемых источников данных есть дополнительная проблема; время привязки к данным, поскольку данные могут измениться между временем создания spliterator и временем выполнения потоковой обработки. В идеале spliterator для потока должен сообщать о характеристике IMMUTABLE или CONCURRENT; в противном случае он должен быть поздне-связываемым. Если источник не может напрямую предоставить рекомендуемый spliterator, он может косвенно предоставить spliterator, используя Supplier, и создать поток с помощью версий, принимающих Supplier, stream(). Spliterator извлекается из поставщика только после начала терминальной операции потокового конвейера.
Эти требования значительно уменьшают область потенциальных помех между мутациями источника потока и выполнением потоковых конвейеров. Потоки, основанные на разделителях с желаемыми характеристиками, или использующие фабричные формы на основе поставщика, защищены от модификаций источника данных до начала терминальной операции (при условии, что параметры поведения операций потока удовлетворяют необходимым критериям отсутствия помех и бессостоятельности). Подробнее см. Непомехи.
- С момента:
- 1.8
| Класс | Описание |
|---|---|
|
BaseStream<T, |
Базовый интерфейс для потоков, которые представляют собой последовательности элементов, поддерживающие последовательные и параллельные агрегатные операции. |
|
Collector<T, |
Операция мутабельного сокращения, которая накапливает входные элементы в изменяемый контейнер результатов, необязательно преобразуя накопленный результат в окончательную форму после обработки всех входных элементов. |
| Collector.Characteristics | Характеристики, указывающие свойства Collector, которые могут быть использованы для оптимизации реализаций сокращения. |
| Collectors | Реализации Collector, которые реализуют различные полезные операции сокращения, такие как накопление элементов в коллекции, подсчет элементов по различным критериям и т.д. |
| DoubleStream | Последовательность примитивных элементов с двойной точностью, поддерживающая последовательные и параллельные агрегатные операции. |
| DoubleStream.Builder | Изменяемый билдер для DoubleStream. |
| DoubleStream.DoubleMapMultiConsumer | Представляет операцию, которая принимает аргумент типа double и DoubleConsumer, и не возвращает результат. |
| IntStream | Последовательность примитивных целых элементов, поддерживающая последовательные и параллельные агрегатные операции. |
| IntStream.Builder | Изменяемый билдер для IntStream. |
| IntStream.IntMapMultiConsumer | Представляет операцию, которая принимает аргумент типа int и IntConsumer, и не возвращает результат. |
| LongStream | Последовательность примитивных элементов типа long, поддерживающая последовательные и параллельные агрегатные операции. |
| LongStream.Builder | Изменяемый билдер для LongStream. |
| LongStream.LongMapMultiConsumer | Представляет операцию, которая принимает аргумент типа long и LongConsumer, и не возвращает результат. |
| Stream<T> | Последовательность элементов, поддерживающая последовательные и параллельные агрегатные операции. |
| Stream.Builder<T> | Изменяемый билдер для Stream. |
| StreamSupport | Вспомогательные методы низкого уровня для создания и управления потоками. |
© 1993, 2021, 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://docs.oracle.com/en/java/javase/17/docs/api/java.base/java/util/stream/package-summary.html