Stream API глубже: collect, reduce, flatMap и ловушки parallel-потоков

Языки3 мин чтения
  • #java
  • #stream-api
  • #functional
  • #parallel
  • #collectors

После 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.stream Javadoc — разделы про spliterators, parallel streams и ассоциативность.
  • «Modern Java in Action» (Urma, Fusco, Mycroft) — подробнее про коллекторы и parallel-потоки с примерами.