Spec-Zone.ru › OpenJDK 21

Пакет java.util.stream

package 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())
               .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() вместо этого снимает всю нагрузку по паралелизации операции сокращения, и библиотека может предоставить эффективную параллельную реализацию без дополнительной синхронизации.

Примеры "гаджетов", показанные ранее, демонстрируют, как сокращение сочетается с другими операциями для замены циклов foreach на объёмные операции. Если 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());
     }
 
Или мы могли бы использовать распараллеливаемую форму 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 на результатах.

Примеры ассоциативных операций включают числовое сложение, 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 или данных, которые могут использоваться реализациями для оптимизации выполнения.

Разделители для изменяемых источников данных сталкиваются с дополнительной проблемой; время привязки к данным, так как данные могут измениться между временем создания разделителя и временем выполнения потоковой обработки. В идеале разделитель для потока должен сообщать о характеристике IMMUTABLE или CONCURRENT; в противном случае он должен быть поздней привязкой. Если источник не может напрямую предоставить рекомендуемый разделитель, он может косвенно предоставить разделитель, используя Supplier, и создать поток с помощью версий, принимающих Supplier stream(). Разделитель извлекается из поставщика только после начала заключительной операции потоковой обработки.

Эти требования значительно сужают область потенциальных конфликтов между изменениями источника потока и выполнением потоковых процессов. Потоки, основанные на разделителях с желаемыми характеристиками или использующих фабричные формы на основе поставщика, устойчивы к изменениям источника данных до начала заключительной операции (при условии, что параметры поведения операций потока соответствуют требуемым критериям отсутствия конфликтов и отсутствия состояния). Подробности см. в разделе Несовместимость.

С момента:
1.8
Пакет Описание
java.util
Содержит фреймворк коллекций, некоторые классы поддержки локализации, загрузчик сервисов, свойства, генерацию случайных чисел, классы для разбора и сканирования строк, кодирование и декодирование Base64, битовый массив и несколько вспомогательных классов.
Класс Описание
BaseStream<T,S extends BaseStream<T,S>>
Базовый интерфейс для потоков, которые представляют собой последовательности элементов, поддерживающие последовательные и параллельные агрегатные операции.
Collector<T,A,R>
Операция изменяемого сокращения, которая накапливает элементы ввода в изменяемый контейнер результата, необязательно преобразуя накопленный результат в окончательное представление после обработки всех элементов ввода.
Collector.Characteristics
Характеристики, указывающие свойства Collector, которые могут быть использованы для оптимизации реализаций сокращения.
Collectors
Реализации Collector, которые реализуют различные полезные операции сокращения, такие как накопление элементов в коллекции, подсчёт элементов по различным критериям и т.д.
DoubleStream
Последовательность примитивных элементов с двойной точностью, поддерживающая последовательные и параллельные агрегатные операции.
DoubleStream.Builder
Изменяемый билдер для DoubleStream.
DoubleStream.DoubleMapMultiConsumer
Представляет операцию, которая принимает аргумент типа double и DoubleConsumer, и не возвращает результат.
IntStream
Последовательность примитивных целочисленных элементов, поддерживающая последовательные и параллельные агрегатные операции.
IntStream.Builder
Изменяемый билдер для IntStream.
IntStream.IntMapMultiConsumer
Представляет операцию, которая принимает аргумент типа int и IntConsumer, и не возвращает результат.
LongStream
Последовательность примитивных длинных целочисленных элементов, поддерживающая последовательные и параллельные агрегатные операции.
LongStream.Builder
Изменяемый билдер для LongStream.
LongStream.LongMapMultiConsumer
Представляет операцию, которая принимает аргумент типа long и LongConsumer, и не возвращает результат.
Stream<T>
Последовательность элементов, поддерживающая последовательные и параллельные агрегатные операции.
Stream.Builder<T>
Изменяемый билдер для Stream.
StreamSupport
Вспомогательные методы низкого уровня для создания и управления потоками.

© 1993, 2023, 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/21/docs/api/java.base/java/util/stream/package-summary.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API