Для чего нужны параллельные стримы java
Перейти к содержимому

Для чего нужны параллельные стримы java

  • автор:

Parallel streams

Стримы бывают последовательными ( sequential ) и параллельными ( parallel ). Последовательные выполняются только в текущем потоке, а вот параллельные используют общий пул ForkJoinPool.commonPool() . При этом элементы разбиваются (если это возможно) на несколько групп и обрабатываются в каждом потоке отдельно. Затем на нужном этапе группы объединяются в одну для предоставления конечного результата.

Чтобы получить параллельный стрим, нужно либо вызвать метод parallelStream() вместо stream(), либо превратить обычный стрим в параллельный, вызвав промежуточный оператор parallel() .

list.parallelStream() .filter(x -> x > 10) .map(x -> x * 2) .collect(Collectors.toList()); 
IntStream.range(0, 10) .parallel() .map(x -> x * 10) .sum(); 

Работа с потоконебезопасными коллекциями, разбиение элементов на части, создание потоков, объединение частей воедино, всё это кроется в реализации Stream API . От нас лишь требуется вызвать нужный метод и проследить, чтобы функции в операторах не зависели от каких-либо внешних факторов, иначе есть риск получить неверный результат или ошибку.

Вот так делать нельзя:

final List ints = new ArrayList<>(); IntStream.range(0, 1000000) .parallel() .forEach(i -> ints.add(i)); System.out.println(ints.size()); 

Это код Шрёдингера. Он может нормально выполниться и показать 1000000, может выполниться и показать 869877, а может и упасть с ошибкой

Exception in thread "main" java.lang.ArrayIndexOutOfBoundsException: 666 at java.util.ArrayList.add(ArrayList.java:777). 

Поэтому разработчики настоятельно просят воздержаться от побочных эффектов в лямбдах, то тут, то там говоря в документации о невмешательстве ( non-interference ).

  • Не используйте параллельные стримы везде, где только можно. Затраты на разбиение элементов, обработку в другом потоке и последующее их слияние порой больше, чем выполнение в одном потоке.
  • При использовании параллельных стримов, убедитесь, что нигде нет блокирующих операций или чего-то, что может помешать обработке элементов.
list.parallelStream() .filter(s -> isFileExists(hash(s))) . 

Как работают параллельные стримы?

Основная цель, ради которой в Java 8 был добавлен Stream API – удобство многопоточной обработки.

Обычный стрим будет выполняться параллельно после вызова промежуточной операции parallel() . Некоторые стримы создаются уже многопоточными, например результат вызова Collection#parallelStream() . Для распараллеливания используется единый общий ForkJoinPool.

Внутри реализации потока его сплиттератор оборачивается в AbstractTask , который и отправляется на выполнение в пул. AbstractTask при выполнении считывает estimateSize сплиттератора и текущую степень параллелизма пула. На основе этих данных он принимает решение, распараллелить ли сплиттератор на два методом trySplit() .

У удобства такого решения есть обратная сторона. Так как пул единый, нагрузка распределяется на всех пользователей параллельных стримов в программе. Если в одном потоке выполняются долгие блокирующие операции, это может ударить по производительности в совершенно не связанном с ним другом потоке.

Если всё же требуется использовать отдельный пул потоков, сам стрим выполняется как задача этого отдельного пула. Подробнее в статье.

Для чего нужны параллельные стримы java

Кроме последовательных потоков Stream API поддерживает параллельные потоки. Распараллеливание потоков позволяет задействовать несколько ядер процессора (если целевая машина многоядерная) и тем самым может повысить производительность и ускорить вычисления. В то же время говорить, что применение параллельных потоков на многоядерных машинах однозначно повысит производительность — не совсем корректно. В каждом конкретном случае надо проверять и тестировать.

Чтобы сделать обычный последовательный поток параллельным, надо вызвать у объекта Stream метод parallel . Кроме того, можно также использовать метод parallelStream() интерфейса Collection для создания параллельного потока из коллекции.

В то же время если рабочая машина не является многоядерной, то поток будет выполняться как последовательный.

Применение параллельных потоков во многих случаях будет аналогично. Например:

import java.util.Optional; import java.util.stream.Stream; public class Program < public static void main(String[] args) < StreamnumbersStream = Stream.of(1, 2, 3, 4, 5, 6); Optional result = numbersStream.parallel().reduce((x,y)-> x*y); System.out.println(result.get()); // 720 > >

Еще один пример:

import java.util.Arrays; import java.util.List; public class Program < public static void main(String[] args) < Listpeople = Arrays.asList("Tom","Bob", "Sam", "Kate", "Tim"); System.out.println("Последовательный поток"); people.stream().filter(p->p.length()==3).forEach(System.out::println); System.out.println("\nПараллельный поток"); people.parallelStream().filter(p->p.length()==3).forEach(System.out::println); > >

В данном случае сначала для списка people создаем поток и выполняем над ним ряд операций в последовательном режиме. В частности, находим в списке строки, длина которых равна 3 и выводим их на консоль. В этом случае все операции с потоком будут производиться над элементами списка в том порядке, в котоом элементы идут в списке.

Затем с помощью метода people.parallelStream() для списка создается параллельный поток. Причем применяются те же операции, однако теперь порядок, в котором над элементами списка будут производиться операции, не детерминирован.

Последовательный поток Tom Bob Sam Tim Параллельный поток Sam Tim Bob Tom

В случае с параллельным потоком вывод недетерминирован и может отличаться.

Однако не все функции можно без ущерба для точности вычисления перенести с последовательных потоков на параллельные. Прежде всего такие функции должны быть без сохранения состояния и ассоциативными, то есть при выполнении слева направо давать тот же результат, что и при выполнении справа налево, как в случае с произведением чисел. Например:

Stream numbersStream = Stream.of(1, 2, 3, 4, 5, 6); Integer result = numbersStream.parallel().reduce(1, (x,y)->x * y); System.out.println(result);

Фактически здесь происходит перемножение чисел. При этом нет разницы между 1 * 2 * 3 * 4 * (5 * 6) или 5 * 6 * 1 * (2 * 3) * 4 . Мы можем расставить скобки любым образом, разместить последовательность чисел в любом порядке, и все равно мы получим один и тот же результат. То есть данная операция является ассоциативной и поэтому может быть распараллелена.

Вопросы производительности в параллельных операциях

Фактически применение параллельных потоков сводится к тому, что данные в потоке будут разделены на части, каждая часть обрабатывается на отдельном ядре процессора, и в конце эти части соединяются, и над ними выполняются финальные операции. Рассмотрим некоторые критерии, которые могут повлиять на производительность в параллельных потоках:

  • Размер данных. Чем больше данных, тем сложнее сначала разделять данные, а потом их соединять.
  • Количество ядер процессора. Теоретически, чем больше ядер в компьютере, тем быстрее программа будет работать. Если на машине одно ядро, нет смысла применять параллельные потоки.
  • Чем проще структура данных, с которой работает поток, тем быстрее будут происходить операции. Например, данные из ArrayList легко использовать, так как структура данной коллекции предполагает последовательность несвязанных данных. А вот коллекция типа LinkedList — не лучший вариант, так как в последовательном списке все элементы связаны с предыдущими/последующими. И такие данные трудно распараллелить.
  • Над данными примитивных типов операции будут производиться быстрее, чем над объектами классов

Упорядоченность в параллельных потоках

Как правило, элементы передаются в поток в том же порядке, в котором они определены в источнике данных. При работе с параллельными потоками система сохраняет порядок следования элементов. Исключение составляет метод forEach() , который может выводить элементы в произвольном порядке. И чтобы сохранить порядок следования, необходимо применять метод forEachOrdered :

phones.parallelStream() .sorted() .forEachOrdered(s->System.out.println(s));

Сохранение порядка в параллельных потоках увеличивает издержки при выполнении. Но если нам порядок не важен, то мы можем отключить его сохранение и тем самым увеличить производительность, использовав метод unordered :

phones.parallelStream() .sorted() .unordered() .forEach(s->System.out.println(s));

Parallel Stream — не панацея или используй с умом (tutorial для начинающих)

Данная статья может быть интересна тем, кто только изучает Stream API, либо набирает практический опыт их использования. В ней раскрывается функционал, плюсы и минусы использования Parallel Stream, но не касаемся использования последовательных Stream API в целом.

Параллельные потоки стали мощной функцией в Java 8 и более поздних версиях, предлагая разработчикам возможность без особых усилий выполнять операции сбора данных параллельно. Используя возможности многопоточности современных компьютеров, параллельные потоки могут значительно повысить производительность вашего кода. В этой статье мы рассмотрим несколько примеров использования параллельных потоков, подчеркнув их преимущества в различных сценариях.

Для чего нужно параллелить потоки:

1. Обработка больших наборов данных

Один из наиболее распространенных вариантов использования параллельных потоков — это работа с большими наборами данных. Допустим, у вас есть список из миллиона записей, и вам нужно выполнить интенсивную вычислительную операцию над каждым элементом. Параллельные потоки могут разделить рабочую нагрузку между несколькими потоками, что значительно сократит время обработки. Например, вы можете использовать параллельные потоки для эффективного выполнения сложных вычислений, фильтрации, сопоставления или группировки операций с большими наборами данных.

2. Операции с интенсивным использованием ЦП

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

3. Улучшенная производительность ввода/вывода

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

4. Сложные преобразования данных

При работе со сложными преобразованиями данных параллельные потоки могут упростить код и улучшить его читаемость. Рассмотрите сценарии, в которых вам нужно применить несколько операций, таких как фильтрация, сопоставление и сокращение, для преобразования набора объектов. Параллельные потоки могут эффективно обрабатывать промежуточные шаги, что приводит к более чистому и лаконичному коду. Это особенно полезно в таких сценариях, как обработка файлов журналов, анализ больших документов XML/JSON или преобразование данных в заданиях пакетной обработки.

5. Сокращение потока

Параллельные потоки являются ценным инструментом при сокращении потоков, например при суммировании, поиске максимума/минимума или накоплении значений. Используя возможности параллельных вычислений, эти операции могут выполняться одновременно, что приводит к значительному повышению производительности. Например, вычисление суммы большого набора чисел с использованием параллельных потоков может дать более быстрые результаты по сравнению с последовательным аналогом.

Примеры использования Parallel Stream

Сам код использования параллельных потоков не сложный и чуть отличается от последовательных Stream:

1) Распараллеливание операции List

List numbers = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10); numbers.parallelStream() .map(n -> n*2) .forEach(System.out::println);

В этом примере показано, как использовать параллельный поток для умножения каждого числа в списке на 2 и вывода результата. Обработка выполняется параллельно, что может повысить производительность для больших наборов данных.

2) Распараллеливание операции сложения

List numbers = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10); int summ = numbers.parallelStream() .reduce(0, Integer::sum);

Параллельный поток используется для вычисления суммы всех чисел в списке. Операция reduce объединяет значения с помощью функции Integer::sum , начиная с начального значения 0.

В чем разница stream().parallel() & parallelStream()?

Методы stream().parallel() и parallelStream() в Java представляют два разных способа создания параллельного потока.

1) stream().parallel() : этот метод используется для преобразования последовательного потока в параллельный поток. Его можно вызвать для любого объекта последовательного потока, чтобы включить параллельную обработку этого потока. Например:

List list = Arrays.asList("a", "b", "c"); Stream parallelStream = list.stream().parallel();

В этом случае метод stream() создает последовательный поток из списка, а затем метод parallel() вызывается для последовательного потока, чтобы преобразовать его в параллельный поток.

2) parallelStream() : этот метод вызывается непосредственно для объекта коллекции для создания параллельного потока. Он возвращает параллельный поток, позволяющий выполнять параллельную обработку элементов коллекции. Например:

List list = Arrays.asList("a", "b", "c"); Stream parallelStream = list.parallelStream();

В этом случае метод parallelStream() вызывается непосредственно в списке для создания параллельного потока.

Оба метода достигают одного и того же результата создания параллельного потока, но основное различие заключается в способе их вызова. Метод stream().parallel() вызывается для последовательного объекта потока, тогда как метод parallelStream() вызывается непосредственно для объекта коллекции.

Parallel Stream & ForkJoinPool

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

Существующая связь между параллельными потоками и инфраструктурой ForkJoinPool заключается в базовой реализации параллельных потоков. Когда вы создаете параллельный поток, он использует ForkJoinPool по умолчанию, предоставленный Java, для параллельного выполнения операций потока. Это означает, что работа по разделению данных и их распределению по нескольким потокам выполняется ForkJoinPool.

ForkJoinPool управляет пулом рабочих потоков и планирует выполнение подзадач параллельного потока этими потоками. Он динамически регулирует количество потоков в зависимости от доступных ядер ЦП и рабочей нагрузки. Таким образом обеспечивается эффективное использование системных ресурсов и повышается общая производительность обработки параллельных потоков.

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

Вместо заключения

Введение ParallelStream в Java произвело революцию в том, как мы выполняем потоковые операции. Благодаря возможности распределять задачи между несколькими потоками ParallelStream предлагает значительный прирост производительности для операций, требующих больших вычислительных ресурсов. Кроме того, его бесшовная интеграция с Stream API устраняет необходимость в ручном управлении потоками, упрощая процесс разработки. Используя ParallelStream, разработчики могут легко извлечь выгоду из преимуществ параллельной обработки, открывая новые уровни эффективности и скорости в своих приложениях.

Однако стоит держать в памяти тот момент, что ParallelStream дает нам «условно бесплатную» многопоточность за счет ресурсов имеющегося ForkJoinPool, что порой может быть просто не оправдано с точки зрения использования.

  • java
  • stream api
  • junior developer
  • стримы
  • многопоточность
  • многопоточное программирование
  • multithreading
  • Программирование
  • Java
  • Учебный процесс в IT

Добавить комментарий

Ваш адрес email не будет опубликован. Обязательные поля помечены *