Spec-Zone.ru › OpenJDK 25

Пакет 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() может повысить производительность некоторых операций с состоянием или терминальных операций. Однако большинство конвейеров потоков, например приведённый выше пример «суммы весов блоков», эффективно распараллеливаются даже при ограничениях на порядок.

Операции свёртки

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

Расширяемость

Реализация Collector, использование фабричного метода java.util.stream.Collector.of(...) или применение предопределённых сборщиков из Collectors позволяет создавать пользовательские повторно используемые терминальные операции.

Реализация Gatherer, использование фабричных методов java.util.stream.Gatherer.of(...) и java.util.stream.Gatherer.ofSequential(...) или применение предопределённых сборщиков из Gatherers позволяет создавать пользовательские повторно используемые промежуточные операции.

Свёртка, параллелизм и упорядоченность

При некоторых сложных операциях свёртки, например при использовании collect(), возвращающего Map, например:
    Map<Buyer, List<Transaction>> salesByBuyer
        = txns.parallelStream()
              .collect(Collectors.groupingBy(Transaction::getBuyer));
, параллельное выполнение может оказаться неэффективным. Это связано с тем, что этап комбинирования (слияние одного Map с другим по ключу) может быть затратным для некоторых реализаций Map.

Однако предположим, что контейнером результата в этой свёртке служит коллекция, допускающая конкурентное изменение, например ConcurrentHashMap. В этом случае параллельные вызовы аккумулятора могут одновременно помещать результаты в один общий контейнер результата, и комбайнеру не придётся объединять отдельные контейнеры. Это может повысить производительность параллельного выполнения. Такая операция называется параллельной свёрткой.

Collector, поддерживающий параллельное сведение, помечается характеристикой Collector.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(Function) — параллельный эквивалент 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. Сплитератор — параллельный аналог Iterator; он описывает (возможно, бесконечную) коллекцию элементов и поддерживает последовательный переход к следующему элементу, обход всех элементов и разделение части входных данных в другой сплитератор, который можно обрабатывать параллельно. На самом низком уровне все потоки управляются сплитератором.

Существует несколько способов реализации сплитератора, и почти все они представляют собой компромисс между простотой реализации и производительностью потоков во время выполнения. Самый простой, но наименее производительный способ создания сплитератора — создать его из итератора с помощью Spliterators.spliteratorUnknownSize(java.util.Iterator, int). Такой сплитератор будет работать, но, скорее всего, обеспечит низкую производительность при параллельной обработке: мы теряем информацию о размере (размере базового набора данных), а также ограничены простым алгоритмом разделения.

Более качественный сплитератор обеспечивает сбалансированное разбиение на части известного размера, точную информацию о размере и ряд других characteristics сплитератора или данных, которые реализации могут использовать для оптимизации выполнения.

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

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

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

Сообщить об ошибке или предложить улучшение
Дополнительные справочные материалы по API и документацию для разработчиков см. в разделе документации Java SE, где приведены более подробные описания для разработчиков, концептуальные обзоры, определения терминов, обходные решения и примеры работающего кода. Другие версии.
Java является товарным знаком или зарегистрированным товарным знаком Oracle и/или её аффилированных лиц в США и других странах.
Авторское право © 1993, 2025, Oracle и/или её аффилированные лица, 500 Oracle Parkway, Redwood Shores, CA 94065 USA.
Все права защищены. Использование регулируется условиями лицензии и политикой распространения документации.

© 1993, 2025, Oracle and/or its affiliates. All rights reserved.
Documentation extracted from Debian's OpenJDK Development Kit package.
Licensed under the GNU General Public License, version 2, with the Classpath Exception.
Various third party code in OpenJDK is licensed under different licenses (see Debian package).
Java and OpenJDK are trademarks or registered trademarks of Oracle and/or its affiliates.
https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/stream/package-summary.html

Spec-Zone.ru

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