Прекратите использовать parallelStream вслепую! | Тренировочный лагерь для авторов, этап 2

Java

предисловие

Многие разработчики хорошо понимают особенности JAVA8, и одна из самых важных функций —Stream, функция, позволяющая создавать сложные запросы к коллекциям декларативным образом. а также,Stream APIОн также предоставляет простой метод для параллельного выполнения. просто добавьparallel()утверждение или использованиеparallelStream()функция. Но если разработчики будут слепо использовать параллельные потоки, это не только не улучшит производительность, но и приведет к фатальным ошибкам.

Пример

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

 public static void main(String[] args) {
        long[] array = new long[]{3, 123, 1, 31, 56, 61, 22};
        long total = Arrays.stream(array)
                .reduce(1, (acc, next) -> acc * next);
        System.out.println(total);
    }

Если полученный результат также необходимо умножить на фиксированное число m, то нам просто нужно изменить код на:

 int total = Arrays.stream(array)
                .reduce(m, (acc, next) -> acc * next);

Если чисел слишком много, не приведет ли последовательное выполнение последовательного потока к неэффективности? Поэтому я попробовал еще разparallel()для выполнения программы.

 public static void main(String[] args) {
        long[] array = new long[]{3, 123, 1, 31, 56, 61, 22};
        long total = Arrays.stream(array)
                .parallel()
                .reduce(1, (acc, next) -> acc * next);
        System.out.println(total);
    }

Я случайно обнаружил, что когдаm=1Когда результаты, полученные последовательным потоком и параллельным потоком, согласуются, но когда m не равно 1, результаты двух не совпадают. например, когдаm=3При , результатом операции последовательного потока является2578991184Результат работы параллельного потока1880084573136. Что вызвало такую ​​ошибку?

ForkJoinPool

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

Stream.reduceПри последовательном выполнении это выглядит так:

未命名文件 (6).png

Алгоритм для параллельных потоков на самом деле очень прост, будем считать, что задача разбивается только на 2 части:

未命名文件 (8).png

Каждый блок умножается на m еще раз, и параллельный поток применяет заданный идентификатор m к каждому блоку задачи. Зная ошибку только сейчас, мы можем решить ее. Мы можем взять 1 для каждого флага m, умножить на 1, не влияя на результат программы, и получить окончательный результат, просто умножив на m:

 public static void main(String[] args) {
        long[] array = new long[]{3, 123, 1, 31, 56, 61, 22};
        long total = Arrays.stream(array)
                .parallel()
                .reduce(1, (acc, next) -> acc * next) * m;
        System.out.println(total);
    }

В этом примере, когда мы используем потоки, на какие мелкие детали мы должны обращать внимание?

Уменьшение должно быть разделяемым

Если вы не уверены, что поток является последовательным потоком (например, он предоставляется в качестве параметра функции), функция сокращенияidentityНе должно влиять на результаты отдельных блоков задач. То есть идентичность функции суммирования должна быть равна 0, а идентичность произведения должна быть равна 1.

Разумное использование параллельных потоков

Не все потоковые операции следует распараллеливать. Напримерmap,flatMapа такжеfilterне имеет состояния, поэтому мы можем использовать подход параллельных потоков. а такжеsort,distinctа такжеlimitЭто не только не принесет улучшения производительности, но может вызвать ошибки.

И эффективность распараллеливания сильно зависит от источника потока.ArrayList,arrayилиIntStream.rangeПоддерживается произвольный доступ, что означает, что их можно легко разделить. ноLinkedListРазложение занимает O(n) времени. а такжеStream.iterateа такжеBufferedReaderТакже старайтесь избегать параллельных потоков, так как все они имеют неизвестную длину в начале, что затрудняет оценку источника разделения.

Пишите модульные тесты

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

Суммировать

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