Stream в Java 8 такой мощный, знаете ли вы его принцип?

задняя часть

API Java 8 добавляет новую абстракцию под названием Stream, которая позволяет обрабатывать данные декларативным способом.

Потоки обеспечивают высокоуровневую абстракцию над операциями над наборами и представлениями Java интуитивно понятным способом, аналогичным запросу данных из базы данных с помощью операторов SQL.

Stream API может значительно повысить производительность Java-программистов, позволяя программистам писать эффективный, чистый и лаконичный код.

В этой статье будет проанализирован принцип реализации Stream.

1. Состав и характеристики Stream

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

  • Элементы — это объекты определенного типа, образующие очередь. Потоки в Java не хранят элементы и не управляют ими, как это делают коллекции, а выполняют вычисления по запросу.

  • Источником потока источника данных может быть коллекция, массив, канал ввода-вывода, генератор и т. д.

  • Операции агрегирования аналогичны операторам SQL, таким как фильтрация, сопоставление, сокращение, поиск, сопоставление, сортировка и т. д.

В отличие от предыдущих операций Collection, у операций Stream есть две основные характеристики:

  • Конвейерная обработка: Промежуточные операции возвращают сам объект потока. Таким образом, несколько операций могут быть объединены в конвейер, как в свободном стиле. Это позволяет оптимизировать такие операции, как оценка лени и короткое замыкание.

  • Внутренняя итерация: ранее обход коллекции выполнялся с помощью Iterator или For-Each, который явно выполняет итерацию за пределами коллекции, что называется внешней итерацией. Stream предоставляет внутренний итеративный метод, реализованный через шаблон посетителя (Visitor).

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

Параллельная работа Stream основана на инфраструктуре Fork/Join (JSR166y), представленной в Java7 для разделения задач и ускорения обработки. Эволюция параллельного API Java в основном выглядит следующим образом:

"

java.lang.Thread в 1.0-1.4

java.util.concurrent в 5.0

Фазеры и многое другое в версии 6.0

Разветвить/объединить платформу в 7.0

Лямбды в 8.0

"

Stream имеет возможности параллельной обработки, и процесс обработки будет разделять и властвовать, то есть большая задача делится на несколько маленьких задач, а это значит, что каждая задача является операцией:

List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9);numbers.parallelStream()       .forEach(out::println);

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

Если это должно быть то же самое, вы можете использовать метод forEachOrdered для выполнения операции завершения:

List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9);numbers.parallelStream()       .forEachOrdered(out::println);

Здесь возникает вопрос, если результат нужно упорядочить, противоречит ли это нашему первоначальному замыслу параллельного выполнения? Да, в этом сценарии очевидно, что нет необходимости использовать параллельные потоки, и его можно выполнить напрямую с последовательными потоками, иначе производительность может быть хуже, потому что все параллельные результаты принудительно сортируются в конце.

Хорошо, давайте сначала представим соответствующие знания интерфейса Stream.

2. Интерфейс BaseStream

Родительским интерфейсом Stream является BaseStream, который представляет собой интерфейс верхнего уровня, реализуемый всеми потоками и определяемый следующим образом:

public interface BaseStream<T, S extends BaseStream<T, S>>        extends AutoCloseable {    Iterator<T> iterator();    Spliterator<T> spliterator();    boolean isParallel();    S sequential();    S parallel();    S unordered();    S onClose(Runnable closeHandler);    void close();}

Среди них T — тип элементов в потоке, S — класс реализации BaseStream, элементы в нем — тоже T и S — тоже он сам:

S extends BaseStream<T, S>

У вас немного кружится голова?

На самом деле это легко понять.Давайте посмотрим на использование S в интерфейсе: например, sequence() и parallel(), эти два метода возвращают экземпляры S, что означает, что они соответственно поддерживают сериализацию текущего потока. Или работать параллельно и возвращать «измененный» объект потока.

"

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

"

Другими словами, S здесь является классом реализации BaseStream, который также является потоком, таким как Stream, IntStream, LongStream и т. д.

3. Потоковый интерфейс

Давайте взглянем на объявление интерфейса Stream:

public interface Stream<T> extends BaseStream<T, Stream<T>>

Это нетрудно понять применительно к приведенному выше объяснению: то есть Stream может продолжать разбиваться на Stream, и мы можем подтвердить это через некоторые его методы:

Stream<T> filter(Predicate<? super T> predicate);<R> Stream<R> map(Function<? super T, ? extends R> mapper);<R> Stream<R> flatMap(Function<? super T, ? extends Stream<? extends R>> mapper);Stream<T> sorted();Stream<T> peek(Consumer<? super T> action);Stream<T> limit(long maxSize);Stream<T> skip(long n);...

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

4. Закройте потоковую операцию

BaseStream реализует интерфейс AutoCloseable, то есть метод close() будет вызываться при закрытии потока. В то же время BaseStream также предоставляет нам метод onClose():

S onClose(Runnable closeHandler);

Когда вызывается интерфейс close() AutoCloseable, будет вызываться метод onClose() объекта потока, но следует отметить несколько моментов:

  • Метод onClose() возвращает сам объект потока, что означает, что объект может вызываться несколько раз.

  • Если вызывается несколько методов onClose(), он срабатывает в том порядке, в котором они вызываются, но если у метода есть исключение, вверх будет выброшено только первое исключение.

  • Исключение, созданное предыдущим методом onClose(), не повлияет на использование последующих методов onClose().

  • Если несколько методов onClose() генерируют исключения, будет отображаться только стек первого исключения, тогда как другие исключения будут сжаты, и будет отображаться только часть информации.

5. Параллельные и последовательные потоки

Интерфейс BaseStream предоставляет два метода, параллельный поток и последовательный поток.Эти два метода могут вызываться произвольно несколько раз или смешиваться, но в конечном итоге превалирует только возвращаемый результат последнего вызова метода.

Обратитесь к описанию метода parallel():

"

Returns an equivalent stream that is parallel. May return

itself, either because the stream was already parallel, or because

the underlying stream state was modified to be parallel.

"

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

В следующем примере преобладает последний вызов parallel(), и сумма, наконец, вычисляется параллельно:

stream.parallel()   .filter(...)   .sequential()   .map(...)   .parallel()   .sum();

6. Человек, стоящий за ParallelStream: ForkJoinPool

Платформа ForkJoin — это новая функция JDK 7. Как и ThreadPoolExecutor, она также реализует интерфейсы Executor и ExecutorService. Он использует «бесконечную очередь» для сохранения задач, которые необходимо выполнить, и количество потоков передается через конструктор.Если желаемое количество потоков не передается в конструктор, количество процессоров, доступных для текущего компьютеру будет установлено значение по умолчанию для количества потоков.

ForkJoinPool в основном используется для решения задач с использованием алгоритма «разделяй и властвуй», типичных приложений, таких как _алгоритм быстрой сортировки_. Дело в том, что ForkJoinPool необходимо использовать относительно небольшое количество потоков для обработки большого количества задач.

Например, для сортировки 10 миллионов данных задача будет разделена на две задачи сортировки по 5 миллионов и задачу слияния для двух наборов данных по 5 миллионов.

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

"

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

"

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

Так в чем же разница в производительности при использовании ThreadPoolExecutor или ForkJoinPool?

Во-первых, с помощью ForkJoinPool можно использовать ограниченное количество потоков для выполнения очень большого количества задач с «отношением родитель-потомок», например, используя 4 потока для выполнения более 2 миллионов задач. При использовании ThreadPoolExecutor завершение невозможно, потому что Thread в ThreadPoolExecutor не может выбрать выполнение подзадач первым.Когда нужно выполнить 2 миллиона задач с отношениями родитель-потомок, также требуется 2 миллиона потоков, что явно невыполнимо.

Принцип кражи работы:

  1. Каждый рабочий поток имеет свою собственную рабочую очередь WorkQueue;

  2. Это двусторонняя очередь из очереди, которая является частным потоком;

  3. Подзадача fork в ForkJoinTask будет помещена в голову очереди рабочего потока, выполняющего задачу, и рабочий поток будет обрабатывать задачи в рабочей очереди в порядке LIFO, то есть в порядке стека;

  4. Чтобы максимизировать использование ЦП, простаивающие потоки будут «воровать» задачи из очередей других потоков для выполнения.

  5. Но он крадет задачи из хвоста рабочей очереди, чтобы уменьшить конкуренцию с потоком, которому принадлежит очередь;

  6. Работа двусторонней очереди: push()/pop() вызывается только в рабочем потоке своего владельца, poll() вызывается, когда другие потоки крадут задачи;

  7. Когда останется только последнее задание, конкуренция все равно будет, что достигается через CAS;

7. Посмотрите на ParallelStream с точки зрения ForkJoinPool

Java 8 добавила общий пул потоков в ForkJoinPool для обработки задач, которые не были явно отправлены в какой-либо пул потоков. Это статический элемент типа ForkJoinPool, который имеет количество потоков по умолчанию, равное количеству ЦП на работающем компьютере.

Автоматическое распараллеливание происходит при вызове нового метода, добавленного в класс Arrays.

Например, параллельная быстрая сортировка, используемая для сортировки массива, используется для параллельного обхода элементов массива. Автоматическое распараллеливание также реализовано в недавно добавленном Stream API Java 8.

Например, следующий код используется для перебора элементов в списке и выполнения требуемой операции:

List<UserInfo> userInfoList =        DaoContainers.getUserInfoDAO().queryAllByList(new UserInfoModel());userInfoList.parallelStream().forEach(RedisUserApi::setUserIdUserInfo);

Операции над элементами в списке выполняются параллельно. Метод forEach создает задачу для операции вычисления каждого элемента, которая обрабатывается commonPool в упомянутом выше ForkJoinPool.

Вышеупомянутая логика параллельных вычислений, конечно, также может быть реализована с помощью ThreadPoolExecutor, но с точки зрения удобочитаемости кода и объема кода использование ForkJoinPool значительно лучше.

Для количества потоков в общем пуле потоков ForkJoinPool обычно можно использовать значение по умолчанию, то есть количество процессоров исполняемого компьютера. Вы также можете настроить количество потоков ForkJoinPool, установив системное свойство: -Djava.util.concurrent .ForkJoinPool.common.parallelism=N (N — количество потоков).

Стоит отметить, что текущий исполняемый поток также будет использоваться для выполнения задачи, поэтому конечное количество потоков равно N+1, а 1 — это текущий основной поток.

Здесь есть проблема: если вы используете _блокирующие операции_ при выполнении параллельных потоков, таких как ввод-вывод, это, скорее всего, вызовет некоторые проблемы:

public static String query(String question) {  List<String> engines = new ArrayList<String>();  engines.add("http://www.google.com/?q=");  engines.add("http://duckduckgo.com/?q=");  engines.add("http://www.bing.com/search?q=");  // get element as soon as it is available  Optional<String> result = engines.stream().parallel().map((base) - {    String url = base + question;    // open connection and fetch the result    return WS.url(url).get();  }).findAny();  return result.get();}

Этот пример типичен, давайте разберем его:

  • Эта операция параллельных потоковых вычислений будет выполняться совместно основным потоком и JVM по умолчанию ForkJoinPool.commonPool().

  • Карта — это метод блокировки, которому необходимо получить доступ к интерфейсу HTTP и получить ответ, поэтому любой рабочий поток будет блокироваться и ждать результата, когда он будет выполняться здесь.

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

  • Текущая реализация ForkJoinPool не рассматривает компенсацию ожидания рабочих потоков, которые заблокированы в ожидании вновь порожденного потока, поэтому в конечном итоге потоки в ForkJoinPool.commonPool() будут находиться в режиме ожидания и блокировать ожидание.

"

Как мы знаем из анализа ситуации выше, выполнение лямбд не является мгновенным, все программы, использующие параллельные потоки, могут стать источником блокирующих программ, а другие части программы не смогут получить доступ к этим воркерам в процессе выполнения ., а это означает, что любая программа, использующая параллельные потоки, станет непредсказуемой и потенциально опасной, когда что-то еще займет общий ForkJoinPool.

"

резюме:

  1. При работе с рекурсивными алгоритмами «разделяй и властвуй» рассмотрите возможность использования ForkJoinPool.

  2. Тщательно установите порог, при котором больше не выполняется разделение задач, этот порог влияет на производительность.

  3. Некоторые функции в Java 8 используют общий пул потоков в ForkJoinPool. В некоторых случаях необходимо настроить количество потоков по умолчанию в пуле потоков.

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

  5. лямбда-выражения не должны мешать друг другу, т. е. избегайте изменения источника данных (поскольку это может привести к проблемам с безопасностью потоков).

  6. Избегайте доступа к состоянию, которое может измениться в течение времени жизни потоковой операции.

8. Производительность параллельных потоков

На производительность инфраструктуры параллельной потоковой передачи влияют следующие факторы:

  • Размер данных: если данные достаточно велики и время обработки каждого конвейера достаточно велико, параллелизм имеет смысл;

  • Структура исходных данных: каждая операция конвейера основана на исходном источнике данных, обычно коллекции, и разделение различных источников данных коллекции потребует определенного объема;

  • Бокс: обработка типов-примитивов выполняется быстрее, чем типов-боксов;

  • Количество ядер: по умолчанию, чем больше ядер, тем больше потоков запускается в базовом пуле потоков fork/join;

  • Накладные расходы на единицу обработки: чем больше времени затрачивается на каждый элемент в потоке, тем более очевидным является улучшение производительности за счет параллельных операций;

Структура исходных данных делится на следующие 3 группы:

  • Хорошая производительность: ArrayList, array или IntStream.range (данные поддерживают случайное чтение и могут быть легко разделены произвольно)

  • Общая производительность: HashSet, TreeSet (данные сложно разложить честно, большинство из них тоже возможно)

  • Низкая производительность: LinkedList (необходимо пройти через цепочку таблиц, трудно разбить полудекомпозицию), stream.Iterate и BufferedReader.Lines (длина неизвестна, сложно разложить)

Примечание: следующая часть раздела выделена из пути: закулисье стримов, спасибо автору _brian gaetz_, писать слишком прозрачно.

9. Модель NQ

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

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

Например, вычисление длины строки требует гораздо меньше усилий, чем вычисление хэша SHA-1 строки. Чем больше работы выполняется для каждого элемента, тем ниже порог «достаточно большой, чтобы воспользоваться преимуществами параллелизма». Точно так же, чем больше у вас данных, тем больше сегментов вы можете разделить, не нарушая порог «слишком маленький».

Простая, но полезная модель параллельной производительности — это модель NQ, где N — количество элементов данных, а Q — объем работы, выполняемой для каждого элемента. Чем больше произведение N*Q, тем больше вероятность, что вы получите параллельное ускорение. Для задач с очень малым Q, таких как суммирование чисел, обычно может потребоваться N > 10 000 для ускорения; по мере увеличения Q размер данных, необходимых для получения ускорения, уменьшается.

Многие препятствия для распараллеливания (такие как разделение затрат, объединение затрат или чувствительность к порядку) могут быть смягчены операциями с более высокой добротностью. Хотя результаты разделения функции LinkedList могут быть плохими, все же возможно получить параллельное ускорение, если у вас достаточно большое значение Q.

10. Последовательность встреч

Порядок встречи относится к тому, является ли порядок, в котором элементы распределяются источником, критическим для вычислений. Некоторые источники (например, наборы и карты на основе хэшей) не имеют значимого порядка встреч. Флаг потока ORDERED указывает, имеет ли поток осмысленный порядок встреч.

Разделитель коллекции JDK установит этот флаг в соответствии со спецификацией коллекции;

Некоторые промежуточные операции могут вводить ORDERED (sorted()) или очищать его (unordered()).

Если поток не соответствует порядку, большинство операций потока должны подчиняться этому порядку. При последовательном выполнении порядок встреч «автоматически сохраняется», потому что элементы обрабатываются естественным образом в том порядке, в котором они встречаются.

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

Но для других операций (промежуточных операций с сохранением состояния, операций завершения, семантика которых связана с порядком встречи, таких как findFirst() или forEachOrdered()), ответственность за соблюдение порядка встречи при параллельном выполнении может быть значительной.

Если существует определенный порядок, порядок не имеет смысла, и порядок конвейера, содержащего последовательную конфиденциальную операцию, можно ускорить, используя операцию unordered() для удаления порядка.

В качестве примера операции, чувствительной к порядку встреч, рассмотрим limit(), которая усекает поток до заданного размера. Реализовать limit() в последовательном выполнении просто: ведите счетчик количества просмотренных элементов, после этого отбрасывайте любые элементы.

Но при параллельном выполнении реализовать limit() намного сложнее, нужно сохранить первые N элементов. Это требование значительно ограничивает возможности использования параллелизма: если ввод разделен на части, вы не будете знать, будет ли результат части включен в окончательный результат, пока не будут завершены все части, предшествующие этой части.

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

Если поток не соответствует порядку, операция limit() может свободно выбирать любые N элементов, что делает выполнение намного более эффективным. Элементы могут быть отправлены вниз по течению, как только они станут известны, без какой-либо буферизации, и единственная координация, которую необходимо выполнить между потоками, — это отправить сигнал, чтобы убедиться, что длина целевого потока не превышена.

Другим менее распространенным примером последовательных затрат является сортировка. Операция sorted() обеспечивает стабильное упорядочение (те же элементы появляются в выходных данных в том же порядке, в котором они были введены во входные данные), если порядок встреч имеет смысл, тогда как для неупорядоченных потоков стабильность (со стоимостью) не требуется.

Аналогичная ситуация и у Different(): если поток имеет порядок встреч, то для нескольких одинаковых входных элементов, Different() должен генерировать первый из них, а для неупорядоченных потоков он может генерировать любой элемент — опять же Гораздо более эффективная параллельная реализация может быть получен.

Аналогичная ситуация возникает при использовании функции collect() для агрегирования. Если операция collect(groupingBy()) выполняется для неупорядоченного потока, элементы, соответствующие любым ключам, должны быть предоставлены нижестоящему сборщику в том порядке, в котором они появляются во входных данных.

Этот порядок обычно не имеет смысла для приложений, и никакой порядок не имеет смысла. В этих случаях может быть лучше выбрать параллельный сборщик (такой как groupingByConcurrent()), который игнорирует порядок встречи и позволяет всем потокам собираться непосредственно в общую параллельную структуру данных (такую ​​как ConcurrentHashMap), вместо того, чтобы каждый поток собирал в свою собственную промежуточную карту, а затем объединяет промежуточную карту (что может быть дорого).

PS: Чтобы эту статью не нашли, вы можете добавить ее в закладки и поставить лайк, чтобы вы могли легко просматривать и находить ее.