Пакет java.util.stream
Классы для поддержки функциональных операций над потоками элементов, таких как преобразования map-reduce над коллекциями. Например:
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())
.collect(Collectors.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, мы могли бы написать очевидную последовательную форму для каждого элемента:
ArrayList<String> strings = new ArrayList<>();
for (T element : stream) {
strings.add(element.toString());
} Или мы могли бы использовать форму collect, пригодную для распараллеливания: 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 на результатах. Примеры ассоциативных операций включают арифметическое сложение, нахождение минимума и максимума, а также конкатенацию строк.
Создание потоков низкого уровня
До сих пор во всех примерах потоков использовались методы, такие какCollection.stream() или Arrays.stream(Object[]) для получения потока. Как реализуются эти методы, генерирующие потоки? Класс StreamSupport содержит ряд методов низкого уровня для создания потока, все использующие ту или иную форму Spliterator. Spliterator — это параллельный аналог Iterator; он описывает (возможно, бесконечную) коллекцию элементов, с поддержкой последовательного продвижения, массового обхода и разделения части входных данных на другой spliterator, который можно обработать параллельно. На самом низком уровне все потоки управляются spliterator.
Существует ряд вариантов реализации spliterator, почти все из которых представляют собой компромисс между простотой реализации и временем выполнения потоков, использующих этот spliterator. Самый простой, но наименее производительный способ создания spliterator — создание spliterator из iterator с помощью Spliterators.spliteratorUnknownSize(java.util.Iterator, int). Хотя такой spliterator будет работать, он, вероятно, обеспечит низкую производительность в параллельном режиме, поскольку мы потеряли информацию о размере (какой размер имеет основной набор данных), а также ограничены упрощенным алгоритмом разделения.
Spliterator более высокого качества будет предоставлять сбалансированные и известные разделения, точную информацию о размере и ряд других characteristics spliterator или данных, которые могут быть использованы реализациями для оптимизации выполнения.
Разделители для изменяемых источников данных сталкиваются с дополнительной проблемой; синхронизации привязки к данным, так как данные могут измениться между временем создания разделителя и временем выполнения потоковой обработки. В идеале, разделитель для потока должен сообщать характеристику IMMUTABLE или CONCURRENT; в противном случае он должен быть поздне-связанным. Если источник не может напрямую предоставить рекомендуемый разделитель, он может косвенно предоставить разделитель, используя Supplier, и создать поток с помощью версий stream(), принимающих Supplier. Разделитель извлекается из поставщика только после начала терминальной операции потоковой обработки.
Эти требования значительно сокращают область возможного вмешательства между изменениями источника потока и выполнением потоковых конвейеров. Потоки, основанные на разделителях с желаемыми характеристиками или использующие фабричные формы на основе поставщика, защищены от модификаций источника данных до начала терминальной операции (при условии, что параметры поведения операциям потока соответствуют необходимым критериям отсутствия вмешательства и бессостоятельности). Дополнительные сведения см. в разделе Отсутствие вмешательства.
- С:
- 1.8
| Интерфейс | Описание |
|---|---|
| BaseStream<T,S extends BaseStream<T,S>> | Базовый интерфейс для потоков, представляющих собой последовательности элементов, поддерживающих последовательные и параллельные агрегатные операции. |
| Collector<T,A,R> | Операция изменяемого сокращения, которая накапливает входные элементы в изменяемый контейнер результата, необязательно преобразуя накопленный результат в окончательное представление после обработки всех входных элементов. |
| DoubleStream | Последовательность примитивных элементов с двойной точностью, поддерживающих последовательные и параллельные агрегатные операции. |
| DoubleStream.Builder | Изменяемый билдер для |
| IntStream | Последовательность примитивных целочисленных элементов, поддерживающих последовательные и параллельные агрегатные операции. |
| IntStream.Builder | Изменяемый билдер для |
| LongStream | Последовательность примитивных элементов с длинной точностью, поддерживающих последовательные и параллельные агрегатные операции. |
| LongStream.Builder | Изменяемый билдер для |
| Stream<T> | Последовательность элементов, поддерживающих последовательные и параллельные агрегатные операции. |
| Stream.Builder<T> | Изменяемый билдер для |
| Класс | Описание |
|---|---|
| Collectors | Реализации |
| StreamSupport | Методы низкого уровня для создания и управления потоками. |
| Перечисление | Описание |
|---|---|
| Collector.Characteristics | Характеристики, указывающие свойства |
© 1993, 2020, 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/11/docs/api/java.base/java/util/stream/package-summary.html