Stream API: продвинутые приёмы
Прогнал одну и ту же сумму двух миллионов чисел. На ArrayList параллельный стрим дал 6683 микросекунды против 15721 у последовательного - вдвое быстрее. На LinkedList тот же переход дал 32833 против 20695: параллельный оказался в полтора раза МЕДЛЕННЕЕ последовательного.
Интервьюеры любят ловить именно на «добавлю parallel и станет быстрее». Проверяют, понимаешь ли ты, какие потоки исполняют конвейер (так называют цепочку операций стрима), как делится источник и когда параллельность вредит.
// Формулировки: «когда parallelStream ускорит, а когда навредит?», «в каком пуле он исполняется?», «зачем нужен IntStream?».
Кто на самом деле исполняет параллельный стрим
Запустил параллельный стрим и записал имена потоков, которые в нём работали. Их оказалось два: ForkJoinPool.commonPool-worker-1 и main. Машина двухъядерная, и общий пул отчитался о параллелизме 1.
Отсюда два факта, которые надо знать. Первый: параллельные стримы исполняются в ОДНОМ общем пуле на всю программу - его называют commonPool, и размер у него по умолчанию на единицу меньше числа ядер. Второй: поток, который позвал терминальную операцию, работает наравне с пулом, поэтому на двух ядрах трудятся два потока, а не один.
Источник делит на куски объект, который называется Spliterator - «делитель». Массив или ArrayList он режет пополам мгновенно по индексам. LinkedList устроен цепочкой ссылок, и чтобы отдать вторую половину, её приходится обойти - отсюда мой замер, где параллельность на LinkedList оказалась хуже обычного обхода.
- commonPool
- общий пул потоков, в котором идут все параллельные стримы программы
- Spliterator
- объект, который умеет делить источник на куски для параллельной обработки
Когда параллельность окупается, а когда нет
Тот же миллион элементов, но работы на каждый я дал заметно больше - двести арифметических операций. ArrayList: 181839 микросекунд последовательно против 117962 параллельно. LinkedList: 182236 против 122964. То есть при тяжёлой работе на элемент разница между источниками почти исчезла - оба выиграли примерно в полтора раза.
Сравни с предыдущей карточкой: там работа была копеечная, и LinkedList от параллельности проиграл. Вывод точнее школьного «LinkedList плохо делится»: качество деления решает ровно тогда, когда работа на элемент мала. А если она мала, параллельность и так почти не нужна.
На маленьких данных всё однозначно. Сто элементов, работы на элемент почти нет: последовательно 30 микросекунд, параллельно 64. Вдвое медленнее - деление, раздача в пул и сборка результата стоят дороже самой работы.
// 1 000 000 элементов, тяжёлая работа на каждый
// ArrayList: 181839 -> 117962 мкс (выигрыш 1,5x)
// LinkedList: 182236 -> 122964 мкс (выигрыш 1,5x)
// 2 000 000 элементов, работа копеечная
// ArrayList: 15721 -> 6683 мкс (выигрыш 2,4x)
// LinkedList: 20695 -> 32833 мкс (ПРОИГРЫШ)
// 100 элементов: 30 -> 64 мкс (ПРОИГРЫШ)- накладные расходы параллельности
- деление источника, раздача в пул, слияние результатов
Два способа сломать параллельный стрим
Первый - побочное действие. Собрал 100 000 чисел параллельным forEach в обычный ArrayList, три раза подряд: 97902, потом 94841, потом ровно 100000. Никаких исключений, просто часть элементов пропала - а на третьей попытке всё сошлось, и ошибку легко счесть выдумкой. Через collect(toList()) все три раза выходило ровно 100000: коллектор устроен так, чтобы правильно сливать куски из разных потоков.
Второй - блокирующее ожидание внутри. Пул общий на всё приложение, поэтому пока твои элементы ждут ответа по сети, встают все остальные параллельные стримы программы. Для задач с ожиданием берут отдельный пул под них или виртуальные потоки - лёгкие потоки, которых можно завести десятки тысяч, потому что на время ожидания они освобождают системный поток. Но не parallelStream.
// И про числа. Сумма миллиона значений через IntStream заняла 540 микросекунд, а через Stream<Integer> - 7967, почти в пятнадцать раз больше. Причина - боксинг: упаковка каждого числа в объект-обёртку Integer. Примитивные стримы (IntStream, LongStream, DoubleStream) обходятся без обёрток и сразу дают sum, average и summaryStatistics; мосты между мирами - mapToInt и boxed.
var out = new ArrayList<Integer>();
ints.parallel().boxed().forEach(out::add);
// три прогона: 97902, 94841, 100000 из 100000 - тихая потеря
ints.parallel().boxed().collect(toList());
// три прогона: 100000, 100000, 100000- боксинг
- упаковка примитивного числа в объект-обёртку, например int в Integer
Как отвечать: «Когда parallelStream поможет, а когда навредит?»
Поможет, когда сходятся три условия: данных много, на каждый элемент приходится ощутимая работа процессора, и конвейер чистый - без блокировок и побочных действий. Навредит в трёх типовых случаях. Маленькая коллекция: у меня сто элементов последовательно считались 30 микросекунд, параллельно 64. Плохо делимый источник при лёгкой работе: два миллиона чисел на ArrayList параллельно дали ускорение вдвое, а на LinkedList - замедление в полтора раза, потому что его приходится обходить, чтобы разделить. И блокирующее ожидание внутри: все параллельные стримы программы идут через один общий пул, я проверял - на двух ядрах это ровно два рабочих потока, включая вызывающий, так что ожидание сети душит всё приложение. Результат собираю только коллектором: параллельный forEach с добавлением в общий список у меня терял по несколько тысяч элементов из ста тысяч, причём не каждый раз.
Почему это сильный ответ: за каждым критерием стоит механика и замер, а не общее правило; названы и случай, когда параллельность выигрывает, и три разных способа проиграть.
На чём валят
- −Ставить parallel «для скорости» на сотне элементов. У меня вышло вдвое медленнее: деление и слияние дороже самой работы.
- −Блокирующее ожидание внутри параллельного стрима. Пул общий на всю программу, и ждать сети в нём означает остановить все параллельные операции.
- −Собирать результат через forEach в общий список. Тихо теряются элементы, и не на каждом прогоне - тесты такое пропускают.
- −Считать, что LinkedList всегда губит параллельность. При тяжёлой работе на элемент он выиграл наравне с ArrayList; разница проявляется на лёгкой.
- −Гонять миллионы чисел через Stream<Integer>. Боксинг обошёлся мне почти в пятнадцать раз против IntStream.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 17, остальные разбираются в тренажёре.
- Когда параллельный стрим (parallelStream) реально ускоряет обработку?A)Parallel быстрее последовательного стрима на данных, поэтому его стоит ставить вездеB)Большой объём, независимые CPU-задачи, ассоциативные операции, без общего состоянияC)На маленьких коллекциях и операциях ввода-вывода — параллелизм максимально раскрывается именно тамD)Когда лямбды меняют общее изменяемое поле — parallel сам синхронизирует к нему доступ безопасно
показать ответ и разбор
+B)Большой объём, независимые CPU-задачи, ассоциативные операции, без общего состояния// разбор: parallelStream использует общий ForkJoinPool и делит данные между ядрами. Выигрыш реален лишь при СОВПАДЕНИИ условий: большой объём, дорогие НЕЗАВИСИМЫЕ вычисления (CPU-bound), операции без состояния и ассоциативные (порядок слияния не важен), источник хорошо делится (массив, ArrayList). Иначе накладные на разбиение/слияние съедают выгоду или он даже медленнее. Опасно: блокирующий I/O (забивает общий пул) и общее изменяемое состояние (гонки — parallel их НЕ защищает). По умолчанию берут последовательный, а parallel — осознанно и с замерами.
- Зачем нужны примитивные стримы IntStream/LongStream/DoubleStream?A)Они работают только с отсортированными числами и автоматически упорядочивают элементы по возрастаниюB)Они позволяют хранить в стриме одновременно int, long и double, автоматически приводя типы между собойC)Избежать боксинга и получить числовые операции (sum, average, range)D)Они предназначены для параллельной обработки и в последовательном режиме недоступны
показать ответ и разбор
+C)Избежать боксинга и получить числовые операции (sum, average, range)// разбор: Stream<Integer> боксит каждый int в Integer (аллокация + разыменование) — на числах это лишняя нагрузка. IntStream/LongStream/DoubleStream работают с примитивами НАПРЯМУЮ (без боксинга) и дают числовые удобства: sum(), average(), max()/min(), summaryStatistics(), а также генераторы range/rangeClosed. Переходы: mapToInt/boxed/asLongStream. Для арифметики по большим наборам это заметно быстрее и экономнее по памяти, поэтому числовые пайплайны пишут на примитивных стримах.
- Как безопасно работать с бесконечным стримом Stream.iterate/generate?A)Бесконечный стрим нужно обходить в отдельном потоке, тогда он сам остановится по завершении метода mainB)Достаточно вызвать sorted() — сортировка заставит бесконечный стрим стать конечным и завершитьсяC)Бесконечный стрим безопасен как есть: терминальная операция count() завершит его за конечное времяD)Обязательно ограничить его (limit или короткозамыкающий терминал)
показать ответ и разбор
+D)Обязательно ограничить его (limit или короткозамыкающий терминал)// разбор: Stream.generate(supplier) и Stream.iterate(seed, f) (двухаргументная форма) порождают БЕСКОНЕЧНЫЙ поток. Работать с ним безопасно только благодаря ленивости + короткому замыканию: limit(n), findFirst, anyMatch — они заставляют стрим выдать конечный результат и остановиться. А вот операции, которым нужно увидеть ВСЕ элементы (sorted, count, collect без limit) — зависнут навсегда. Трёхаргументная Stream.iterate(seed, hasNext, next) (Java 9+) имеет встроенное условие остановки. Правило: бесконечный источник + обязательный ограничитель.
- Для чего предназначен метод peek() и почему на него нельзя полагаться для логики?A)Для отладки/подглядывания; может не выполниться для элементов, не нужных результатуB)Для изменения элементов на месте: peek заменяет map, когда результат преобразования не важен вызывающемуC)Для завершения стрима: peek — терминальная операция, аналогичная forEach, но возвращающая стримD)Для параллельного выполнения побочных действий: peek вызывается для каждого элемента разом
показать ответ и разбор
+A)Для отладки/подглядывания; может не выполниться для элементов, не нужных результату// разбор: peek(consumer) — ПРОМЕЖУТОЧНАЯ операция, задуманная для отладки: «подсмотреть» элементы, проходящие через конвейер (лог, брейкпоинт), не меняя их. Полагаться на неё для бизнес-логики или побочных эффектов НЕЛЬЗЯ: из-за ленивости и оптимизаций (например, count() может не проходить элементы, а findFirst — оборвать) peek может вызваться не для всех элементов или не вызваться вовсе. Для гарантированного действия над каждым элементом есть терминальный forEach. peek — только диагностика.
- Что означают три аргумента reduce(identity, accumulator, combiner)?A)Минимум, максимум и шаг диапазона, в пределах которого выполняется свёртка элементов стрима по порядкуB)Начальное значение, как присоединять элемент, как объединять частичные результаты (для parallel)C)Три альтернативных функции свёртки, из которых стрим в рантайме выбирает самую быструю под данныеD)Начальный, промежуточный и финальный коллекторы, применяемые к стриму строго один за другим по очереди
показать ответ и разбор
+B)Начальное значение, как присоединять элемент, как объединять частичные результаты (для parallel)// разбор: Трёхаргументный reduce используют, когда тип результата ОТЛИЧАЕТСЯ от типа элементов и/или для параллельной свёртки. identity — нейтральное стартовое значение (0 для суммы, "" для строк); accumulator(partial, element) — как присоединить очередной элемент к промежуточному результату; combiner(partial1, partial2) — как СЛИТЬ два промежуточных результата, посчитанных разными потоками в параллельном режиме. Требования: identity нейтрален, операции ассоциативны и согласованы — иначе параллельный результат будет неверным. Для изменяемых контейнеров обычно берут collect, а не такой reduce.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.