Stream API глубже: collect, reduce, flatMap и ловушки parallel-потоков
После filter/map/collect(toList()) многим кажется, что Stream API освоен. На деле это верхушка — а под ней collect с готовыми коллекторами, свёртка reduce, flatMap, короткое замыкание и отдельная история с parallelStream, на которой чаще обжигаются, чем выигрывают. Разберём по порядку.
collect против reduce — в чём разница
Обе операции «сворачивают» поток в одно значение, но делают это принципиально по-разному.
reduce комбинирует элементы попарно через бинарную функцию. Каждый шаг создаёт новое значение, и accumulator должен быть без побочных эффектов и ассоциативным — иначе результат в параллельном потоке будет непредсказуемым.
// Сумма через reduce: создаёт промежуточные значения
int sum = numbers.stream().reduce(0, Integer::sum);
collect работает через изменяемый контейнер (например, ArrayList или HashMap) — он мутирует один и тот же объект, добавляя в него элементы. Это эффективнее по памяти: нет потока промежуточных значений.
// Группировка по первой букве — reduce так не сделает красиво
Map<Character, List<String>> byLetter =
words.stream().collect(Collectors.groupingBy(w -> w.charAt(0)));
Правило простое: надо собрать в коллекцию или сгруппировать — collect. Надо получить одно скалярное значение (сумма, максимум, конкатенация) — reduce. Не наоборот.
Collectors, которые экономят время
groupingBy — самый частый, но у него есть вторая форма с downstream-коллектором:
// Сколько слов на каждую букву
Map<Character, Long> counts =
words.stream().collect(
Collectors.groupingBy(w -> w.charAt(0), Collectors.counting()));
// Средняя длина по группе
Map<Character, Double> avgLen =
words.stream().collect(
Collectors.groupingBy(
w -> w.charAt(0),
Collectors.averagingInt(String::length)));
partitioningBy — частный случай groupingBy с булевым ключом, чуть быстрее и читается лучше для фильтра-на-две-кучи. joining склеивает строки с разделителем. teeing (с Java 12) считает несколько агрегатов за один проход — например, минимум и максимум одновременно, без двух проходов по потоку.
flatMap против map — частая путаница
map превращает каждый элемент в один новый элемент. flatMap — в поток элементов, который потом «сплющивается» в общий поток. Классический кейс — развернуть вложенную структуру:
record Order(int id, List<String> items) {}
// Хотим все товары из всех заказов одним списком
List<String> allItems = orders.stream()
.flatMap(o -> o.items().stream()) // каждый заказ → поток товаров
.distinct()
.toList();
Запомнить: если map возвращает Stream<Stream<X>> и вы честно крутите forEach дважды — нужен flatMap.
Короткое замыкание
Некоторые операции не требуют прохода по всему потоку:
boolean hasAdmin = users.stream().anyMatch(u -> u.role() == Role.ADMIN);
Optional<User> firstActive = users.stream().filter(User::active).findFirst();
Как только anyMatch нашёл совпадение, он останавливается. findFirst тоже. Это важно для производительности на больших потоках — но работает только с &&-семантикой (allMatch, anyMatch, noneMatch, findFirst, findAny). collect и reduce всегда проходят до конца.
parallelStream — где обжигаются
«Поставил parallel() и стало быстрее» — это почти никогда не правда. Причины:
1. Общий ForkJoinPool. Параллельные потоки используют общий ForkJoinPool.commonPool() размером с число ядер минус один. Запустили параллельный поток — заняли слот, который делят с другими задачами в JVM, включая чужие. Два параллельных потока одновременно дерутся за один пул.
2. Разбиение (splitting). Чтобы распараллелить, поток должен уметь хорошо разбиваться на части. ArrayList и массивы разбиваются мгновенно, а LinkedList или поток из BufferedReader.lines() — плохо, потому что splitter должен дойти до середины линейно.
3. Затраты на координацию. Разбить, раздать воркерам, собрать результаты обратно — это не бесплатно. На маленьком потоке (сотни элементов) overhead съест весь выигрыш.
4. I/O в потоке — категорически нельзя. Если внутри лямбды идёт запрос к БД, параллельный поток не ускорит его в N раз — он просто забьёт общий пул блокированными задачами и затормозит всю JVM. Для I/O используют виртуальные потоки (см. отдельную статью), а не parallel-стримы.
Где parallelStream честно помогает: CPU-bound обработка больших объёмов данных (миллионы+ элементов) в ArrayList/массиве, без I/O, с предсказуемой задачей в лямбде. Числа Фибоначчи, обработка изображений, агрегация больших логов.
Ещё пара классических граблей
Автоупаковка примитивов. Stream<Integer> при сложении миллионов элементов создаёт миллионы Integer-объектов. Для примитивов есть IntStream, LongStream, DoubleStream — без упаковки:
// Плохо: миллион упаковок/распаковок
int sum = list.stream().mapToInt(Integer::intValue).sum(); // если list — List<Integer>
// Хорошо для массива:
int sum = IntStream.of(array).sum();
Side effects в лямбдах. Изменение внешней коллекции внутри forEach или map (list.add(...)) — работает, но ломает предсказуемость в параллельном потоке и противоречит концепции. Правильно — собирать через collect.
Нарушение ассоциативности в reduce. Если операция не ассоциативна ((a+b)+c != a+(b+c) — для чисел с плавающей точкой это уже не всегда так), результат параллельной свёртки зависит от порядка и может плавать.
Короткое summary
Stream API — это не «фильтр-мап-коллект», а довольно богатый набор инструментов. collect с готовыми коллекторами покрывает большинство повседневных задач группировки и агрегации; reduce — для скалярных свёрток; flatMap разворачивает вложенные структуры; короткое замыкание экономит проходы. А parallelStream — узкий инструмент для CPU-bound работы над большими данными, не панацея и почти всегда вредная для I/O.
Что почитать
- Oracle Java Tutorial: Aggregate Operations — официальный разбор
collect,reduce, коллекторов. - Package
java.util.streamJavadoc — разделы про spliterators, parallel streams и ассоциативность. - «Modern Java in Action» (Urma, Fusco, Mycroft) — подробнее про коллекторы и parallel-потоки с примерами.