сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Stream API

Stream API: продвинутые приёмы

parallelStream: кнопка «быстрее», которая работает не всегда

Прогнал одну и ту же сумму двух миллионов чисел. На 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, остальные разбираются в тренажёре.

  1. #streams_advanced1 / 5
    Когда параллельный стрим (parallelStream) реально ускоряет обработку?
    A)Parallel быстрее последовательного стрима на данных, поэтому его стоит ставить везде
    B)Большой объём, независимые CPU-задачи, ассоциативные операции, без общего состояния
    C)На маленьких коллекциях и операциях ввода-вывода — параллелизм максимально раскрывается именно там
    D)Когда лямбды меняют общее изменяемое поле — parallel сам синхронизирует к нему доступ безопасно
    показать ответ и разбор
    +B)Большой объём, независимые CPU-задачи, ассоциативные операции, без общего состояния

    // разбор: parallelStream использует общий ForkJoinPool и делит данные между ядрами. Выигрыш реален лишь при СОВПАДЕНИИ условий: большой объём, дорогие НЕЗАВИСИМЫЕ вычисления (CPU-bound), операции без состояния и ассоциативные (порядок слияния не важен), источник хорошо делится (массив, ArrayList). Иначе накладные на разбиение/слияние съедают выгоду или он даже медленнее. Опасно: блокирующий I/O (забивает общий пул) и общее изменяемое состояние (гонки — parallel их НЕ защищает). По умолчанию берут последовательный, а parallel — осознанно и с замерами.

  2. #streams_advanced2 / 5
    Зачем нужны примитивные стримы 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. Для арифметики по большим наборам это заметно быстрее и экономнее по памяти, поэтому числовые пайплайны пишут на примитивных стримах.

  3. #streams_advanced3 / 5
    Как безопасно работать с бесконечным стримом Stream.iterate/generate?
    A)Бесконечный стрим нужно обходить в отдельном потоке, тогда он сам остановится по завершении метода main
    B)Достаточно вызвать 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+) имеет встроенное условие остановки. Правило: бесконечный источник + обязательный ограничитель.

  4. #streams_advanced4 / 5
    Для чего предназначен метод peek() и почему на него нельзя полагаться для логики?
    A)Для отладки/подглядывания; может не выполниться для элементов, не нужных результату
    B)Для изменения элементов на месте: peek заменяет map, когда результат преобразования не важен вызывающему
    C)Для завершения стрима: peek — терминальная операция, аналогичная forEach, но возвращающая стрим
    D)Для параллельного выполнения побочных действий: peek вызывается для каждого элемента разом
    показать ответ и разбор
    +A)Для отладки/подглядывания; может не выполниться для элементов, не нужных результату

    // разбор: peek(consumer) — ПРОМЕЖУТОЧНАЯ операция, задуманная для отладки: «подсмотреть» элементы, проходящие через конвейер (лог, брейкпоинт), не меняя их. Полагаться на неё для бизнес-логики или побочных эффектов НЕЛЬЗЯ: из-за ленивости и оптимизаций (например, count() может не проходить элементы, а findFirst — оборвать) peek может вызваться не для всех элементов или не вызваться вовсе. Для гарантированного действия над каждым элементом есть терминальный forEach. peek — только диагностика.

  5. #streams_advanced5 / 5
    Что означают три аргумента reduce(identity, accumulator, combiner)?
    A)Минимум, максимум и шаг диапазона, в пределах которого выполняется свёртка элементов стрима по порядку
    B)Начальное значение, как присоединять элемент, как объединять частичные результаты (для parallel)
    C)Три альтернативных функции свёртки, из которых стрим в рантайме выбирает самую быструю под данные
    D)Начальный, промежуточный и финальный коллекторы, применяемые к стриму строго один за другим по очереди
    показать ответ и разбор
    +B)Начальное значение, как присоединять элемент, как объединять частичные результаты (для parallel)

    // разбор: Трёхаргументный reduce используют, когда тип результата ОТЛИЧАЕТСЯ от типа элементов и/или для параллельной свёртки. identity — нейтральное стартовое значение (0 для суммы, "" для строк); accumulator(partial, element) — как присоединить очередной элемент к промежуточному результату; combiner(partial1, partial2) — как СЛИТЬ два промежуточных результата, посчитанных разными потоками в параллельном режиме. Требования: identity нейтрален, операции ассоциативны и согласованы — иначе параллельный результат будет неверным. Для изменяемых контейнеров обычно берут collect, а не такой reduce.

дальше

Теорию прочитали. Навык ставится повторением

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