Java Stream API и функциональный стиль: лямбды, коллекторы, ленивость
0%

Stream API и функциональный стиль: лямбды, коллекторы, ленивость

Stream API и функциональный стиль: лямбды, коллекторы, ленивость

Java 8 принесла лямбды и Stream API, и с тех пор в командах сложилось два лагеря. Один пишет через стримы всё, включая трёхэлементные списки и мутацию внешнего состояния в forEach. Другой не пишет через стримы ничего, потому что «в проде они медленные и их не отладить». Оба правы ровно наполовину, и разобраться, где именно, можно только поняв модель исполнения: что такое лямбда в байткоде, что физически происходит при .filter(...).map(...), кто кого тянет — источник конвейер или конвейер источник, и почему parallel() иногда даёт восьмикратное ускорение, а иногда двукратное замедление и неправильный ответ.

Часть 1. Лямбда — это не анонимный класс

Задача, которую решали

До Java 8 передать поведение как значение можно было только объектом — анонимный класс на пять строк ради одного выражения:

Collections.sort(orders, new Comparator<Order>() {
    @Override
    public int compare(Order a, Order b) {
        return a.amount().compareTo(b.amount());
    }
});

Проблема не только в многословности. Каждый анонимный класс — это отдельный class-файл в jar-е (OrderService$1.class), который загрузчик обязан прочитать, верифицировать и разместить в метаспейсе. В крупном приложении их копились тысячи. Поэтому OpenJDK сознательно отказалась от «сахара над анонимным классом»; разбор решения — в Translation of Lambda Expressions Брайана Гоетца.

Что генерирует компилятор

Тело лямбды javac выносит в синтетический статический метод того же класса, а на месте самой лямбды ставит инструкцию invokedynamic — единственную инструкцию JVM, чья семантика определяется не спецификацией, а кодом на Java (bootstrap-методом).

import java.util.function.Predicate;

public class LambdaShape {

    static Predicate<String> nonCapturing() {
        return s -> s.isEmpty();               // ничего не захватывает
    }

    static Predicate<String> capturing(int max) {
        return s -> s.length() > max;          // захватывает max
    }

    public static void main(String[] args) {
        // одинаковое тело, но два РАЗНЫХ места в коде -> два разных класса
        Predicate<String> a = s -> s.isEmpty();
        Predicate<String> b = s -> s.isEmpty();
        System.out.println(a == b);                              // false

        // одно и то же место в коде, захвата нет -> HotSpot переиспользует экземпляр
        System.out.println(nonCapturing() == nonCapturing());    // true

        // захват есть -> на каждый вызов новый объект
        System.out.println(capturing(5) == capturing(5));        // false
    }
}
false
true
false

Поведение nonCapturing() == nonCapturing() -> true не гарантировано спецификацией (JLS §15.27.4 явно оставляет это на усмотрение реализации), но так работает HotSpot, и это важное практическое свойство: лямбда без захвата бесплатна по аллокациям.

Посмотрим в байткод — javap -c -p LambdaShape.class:

static java.util.function.Predicate<java.lang.String> nonCapturing();
  Code:
     0: invokedynamic #7,  0    // InvokeDynamic #0:test:()Ljava/util/function/Predicate;
     5: areturn

static java.util.function.Predicate<java.lang.String> capturing(int);
  Code:
     0: iload_0
     1: invokedynamic #13, 0    // InvokeDynamic #1:test:(I)Ljava/util/function/Predicate;
     6: areturn

private static boolean lambda$capturing$1(int, java.lang.String);   // тело лямбды

BootstrapMethods:
  0: REF_invokeStatic java/lang/invoke/LambdaMetafactory.metafactory:(...)

Захваченные значения — аргументы фабрики ((I)Ljava/util/function/Predicate;). При первом исполнении invokedynamic вызывается LambdaMetafactory.metafactory, которая на лету генерирует hidden class, реализующий целевой интерфейс; дальше call site залинкован и стоит ровно как обычный вызов.

Вывод про старт приложения. Генерация hidden-класса происходит при первом проходе через каждый call site; в Spring-приложении с тысячами лямбд это сотни миллисекунд старта. Лечится AppCDS (-XX:ArchiveClassesAtExit=app.jsa архивирует прелинкованные lambda-классы), а в JDK 24+ — AOT-кэшем (JEP 483). Подробнее про загрузку классов и метаспейс — в JVM изнутри.

Захват и effectively final

Лямбда захватывает значение, а не переменную. Отсюда требование: захватываемая локальная переменная должна быть final или effectively final (JLS §4.12.4).

int counter = 0;
list.forEach(x -> counter++);   // не компилируется: counter не effectively final

Это не каприз компилятора, а следствие того, что лямбда может пережить стековый фрейм с этой переменной. Обход через массив-на-один-элемент (int[] counter = {0}; ... counter[0]++) компилируется — и молча ломается при parallel(). Нужен счётчик — берите LongAdder или, почти всегда правильнее, терминальную операцию count().

Отдельно: поле объекта захватывать можно всегда, потому что на самом деле захватывается this. Это тонкий источник утечек — лямбда в долгоживущем реестре колбэков удерживает весь внешний объект.

Функциональные интерфейсы и ссылки на методы

Функциональный интерфейс — интерфейс с ровно одним абстрактным методом. @FunctionalInterface не обязательна, но включает проверку компилятором и потому обязательна по кодстайлу. В java.util.function их 43 штуки; на практике нужен десяток.

Интерфейс Сигнатура Типичное применение
Function<T,R> T -> R map
BiFunction<T,U,R> (T,U) -> R Map.merge, teeing
Predicate<T> T -> boolean filter, спецификации
Supplier<T> () -> T ленивое создание, orElseGet
Consumer<T> T -> void forEach, побочные эффекты
UnaryOperator<T> T -> T List.replaceAll
BinaryOperator<T> (T,T) -> T reduce, Map.merge
IntFunction, ToIntFunction, IntPredicate без боксинга примитивные конвейеры

Примитивные варианты существуют не для красоты: Function<Integer,Integer> на каждом вызове боксит вход и выход. В горячем цикле это разница на порядок.

Ссылки на методы бывают четырёх видов, и путаница между вторым и третьим — классический источник непонятных ошибок компиляции:

Function<String, Integer> f1 = Integer::parseInt;   // 1. статический:  (a) -> Type.m(a)
Consumer<String> f2 = log::info;                    // 2. конкретный объект: log ЗАХВАТЫВАЕТСЯ
Function<String, Integer> f3 = String::length;      // 3. произвольный объект типа:
                                                    //    (obj, a) -> obj.m(a)
Supplier<ArrayList<String>> f4 = ArrayList::new;    // 4. конструктор

Форма 3 — та самая, из-за которой String::length подходит под Function<String,Integer>, а String::isEmpty — под Predicate<String>: компилятор «сдвигает» получателя в аргументы.

Часть 2. Стрим — это описание конвейера

Что стрим не является

Главное заблуждение: «стрим — это ленивая коллекция». Нет. Стрим:

  • не хранит данные — он знает только, откуда их брать;
  • не мутирует источник — операции возвращают новый стрим;
  • одноразов — повторный терминальный вызов даёт IllegalStateException;
  • может быть бесконечным;
  • ленив — до терминальной операции не выполняется ничего.

Последнее проверяется в три строки:

List<String> names = List.of("anna", "bob", "carol", "dave");

Optional<String> first = names.stream()
        .peek(n -> System.out.println("смотрим: " + n))
        .filter(n -> n.length() > 3)
        .map(n -> {
            System.out.println("в верхний регистр: " + n);
            return n.toUpperCase();
        })
        .findFirst();

System.out.println(first.orElse("нет"));
смотрим: anna
в верхний регистр: anna
ANNA

Ни bob, ни carol, ни dave даже не были прочитаны. Элементы идут через конвейер по одному целиком, а не «слой за слоем». Это свойство называется fusion — слияние операций в единую цепочку Sink; именно оно делает стримы конкурентоспособными с циклом.

Материализация промежуточных коллекций против слияния операций в конвейере Stream

Анатомия конвейера

Три вида операций:

  1. Промежуточные statelessmap, filter, flatMap, mapMulti, peek, mapToInt. Обрабатывают элемент, не помня о других. Сливаются в один проход.
  2. Промежуточные statefulsorted, distinct, limit, skip, takeWhile, dropWhile. Требуют буфера или счётчика. sorted и distinct — полноценные барьеры: до них конвейер обязан прочитать весь вход.
  3. Терминальные — запускают выполнение. Часть из них короткозамкнутые: findFirst, findAny, anyMatch, allMatch, noneMatch, а из промежуточных — limit и takeWhile.

Отсюда правило: стрим нельзя сохранить в поле и переиспользовать. Если нужно «переиспользовать конвейер» — храните Supplier<Stream<T>>:

Supplier<Stream<Order>> active = () -> orders.stream().filter(Order::isActive);
long count = active.get().count();
BigDecimal total = active.get().map(Order::amount).reduce(BigDecimal.ZERO, BigDecimal::add);

Spliterator: контракт источника

Под каждым стримом лежит Spliterator<T> — «итератор, умеющий делиться». Его четыре метода определяют вообще всё поведение конвейера:

  • tryAdvance(Consumer) — обработать один элемент (pull);
  • forEachRemaining(Consumer) — обработать все (быстрый путь для последовательного случая);
  • trySplit() — отдать половину, оставив себе вторую (основа параллелизма);
  • characteristics() + estimateSize() — что известно об источнике.

Характеристики — не документация, а входные данные для оптимизаций:

Флаг Что даёт конвейеру
SIZED / SUBSIZED точный размер: предвыделение массива в toArray, count() без обхода
ORDERED обязанность сохранять порядок; unordered() снимает её и ускоряет distinct/limit в параллели
DISTINCT distinct() вырождается в no-op
SORTED sorted() с тем же компаратором вырождается в no-op
IMMUTABLE / CONCURRENT источник можно не защищать от конкурентной модификации
NONNULL не нужны проверки на null

Практическое следствие, на которое напарываются все:

long n = Stream.of(1, 2, 3).peek(System.out::println).count();
// Java 8:  печатает 1 2 3, n == 3
// Java 9+: не печатает НИЧЕГО,  n == 3

Начиная с Java 9 count() пропускает выполнение конвейера, если размер известен (SIZED) и ни одна операция его не меняет. peek объявлен как операция «для отладки» и на размер не влияет, поэтому его просто выбрасывают. Никогда не используйте peek для полезной работы — это документированная ловушка.

Своя реализация Spliterator — способ подружить ленивый источник со Stream API. Классический кейс — постраничный HTTP-API, который не хочется грузить целиком:

/** Ленивый стрим над постраничным API: страница читается, только когда исчерпана предыдущая. */
static <T> Stream<T> paged(IntFunction<List<T>> fetchPage) {
    var spliterator = new Spliterators.AbstractSpliterator<T>(
            Long.MAX_VALUE,                                    // размер неизвестен
            Spliterator.ORDERED | Spliterator.NONNULL) {

        private int page = 0;
        private Iterator<T> current = Collections.emptyIterator();

        @Override
        public boolean tryAdvance(Consumer<? super T> action) {
            while (!current.hasNext()) {
                List<T> next = fetchPage.apply(page++);         // сетевой вызов
                if (next.isEmpty()) return false;               // страницы кончились
                current = next.iterator();
            }
            action.accept(current.next());
            return true;
        }
    };
    return StreamSupport.stream(spliterator, false);            // false = последовательный
}

// Найдёт первого подходящего и остановится — лишние страницы не запросит.
Optional<User> admin = paged(page -> api.users(page))
        .filter(u -> u.role() == Role.ADMIN)
        .findFirst();

AbstractSpliterator сам реализует наивный trySplit через буферизацию, так что даже такой источник можно параллелить — правда, плохо. trySplit за O(1) умеют только структуры с произвольным доступом.

Часть 3. Операции: что выбирать и почему

map, flatMap, mapMulti

flatMap — «развернуть и склеить»: каждый элемент превращается в стрим, стримы конкатенируются.

record Order(String id, List<Item> items) {}

// все позиции всех заказов одним стримом
List<Item> allItems = orders.stream()
        .flatMap(o -> o.items().stream())
        .toList();

// то же самое, но без аллокации стрима на каждый заказ (Java 16+)
List<Item> same = orders.stream()
        .<Item>mapMulti((order, downstream) -> order.items().forEach(downstream))
        .toList();

Три тонкости, которые стоят денег: flatMap создаёт объект-стрим на каждый элемент (для миллионов элементов — миллионы короткоживущих объектов; mapMulti этого не делает); до Java 10 flatMap ломал короткое замыкание — findFirst после него всё равно вычислял подстрим целиком (JDK-8075939); и flatMap закрывает каждый подстрим после потребления, что важно, если подстримы открывают файлы.

Optional как стрим на 0 или 1 элемент

// вместо filter(Optional::isPresent).map(Optional::get)
List<Config> configs = names.stream()
        .map(registry::find)          // Stream<Optional<Config>>
        .flatMap(Optional::stream)    // Java 9+: пустой Optional просто исчезает
        .toList();

Про сам Optional и его границы применимости (не поле, не параметр, не коллекция) подробно — в идиоматике и обработке ошибок.

Сортировка и компараторы

List<Order> sorted = orders.stream()
        .sorted(Comparator.comparing(Order::customer)
                          .thenComparing(Order::date, Comparator.reverseOrder())
                          .thenComparingInt(Order::priority))   // без боксинга
        .toList();

Три грабли подряд: Comparator.comparing(Order::amount).reversed() разворачивает всю цепочку, а не последний ключ (для одного поля передавайте Comparator.reverseOrder() вторым аргументом, как выше); comparing с примитивным extractor боксит — есть comparingInt/comparingLong/comparingDouble; sorted() и distinct() на бесконечном стриме зависают навсегда без единой ошибки.

Сложность: sorted() — это Arrays.sort (TimSort) поверх материализованного массива, то есть O(N log N) по времени и дополнительно O(N) по памяти. Барьер отменяет всю экономию на fusion — это осознанная плата.

toList и его родственники

List<String> a = stream.collect(Collectors.toList());              // изменяемый (не гарантировано!)
List<String> b = stream.collect(Collectors.toUnmodifiableList());  // Java 10+, NPE на null
List<String> c = stream.toList();                                  // Java 16+, неизменяемый, null ок

По умолчанию берите stream.toList(). Collectors.toList() нужен, только когда результат действительно надо потом менять.

Часть 4. Коллекторы

Устройство

reduce умеет сворачивать в иммутабельное значение. Для сворачивания в мутабельный контейнер (список, map, StringBuilder) есть collect и интерфейс Collector:

supplier создаёт пустой контейнер (у каждого потока свой), accumulator добавляет элемент, combiner сливает два контейнера и вызывается только в параллельном режиме, finisher превращает контейнер в результат, characteristics разрешает оптимизации.

Свой коллектор: топ-N без полной сортировки

Наивное sorted().limit(n) сортирует всё: O(N log N) времени и O(N) памяти. Коллектор на куче даёт O(N log n) времени и O(n) памяти:

/** Собирает n наибольших элементов по компаратору. O(N log n) время, O(n) память. */
static <T> Collector<T, ?, List<T>> topN(int n, Comparator<? super T> cmp) {
    return Collector.of(
            () -> new PriorityQueue<T>(cmp),          // min-heap: наименьший наверху
            (heap, item) -> {                          // accumulator
                heap.offer(item);
                if (heap.size() > n) heap.poll();      // выкидываем наименьший из отобранных
            },
            (left, right) -> {                         // combiner (только для parallel)
                for (T item : right) {
                    left.offer(item);
                    if (left.size() > n) left.poll();
                }
                return left;
            },
            heap -> heap.stream().sorted(cmp.reversed()).toList(),   // finisher
            Collector.Characteristics.UNORDERED);
}

// применение
List<Order> top10 = orders.stream()
        .collect(topN(10, Comparator.comparing(Order::amount)));

Про устройство PriorityQueue — в статье о кучах и очередях с приоритетом.

Готовые коллекторы, которые реально нужны

record Order(String customer, String region, BigDecimal amount, LocalDate date) {}

// 1. Группировка с downstream-коллектором и заданным типом map
Map<String, BigDecimal> revenueByRegion = orders.stream()
        .collect(Collectors.groupingBy(
                Order::region,
                TreeMap::new,                                     // упорядоченный результат
                Collectors.reducing(BigDecimal.ZERO, Order::amount, BigDecimal::add)));

// 2. Двухуровневая группировка
Map<String, Map<Month, Long>> counts = orders.stream()
        .collect(Collectors.groupingBy(
                Order::region,
                Collectors.groupingBy(o -> o.date().getMonth(), Collectors.counting())));

// 3. filtering vs filter: filtering сохраняет пустые группы (Java 9+)
Map<String, List<Order>> bigByRegion = orders.stream()
        .collect(Collectors.groupingBy(
                Order::region,
                Collectors.filtering(o -> o.amount().compareTo(THRESHOLD) > 0, Collectors.toList())));

// 4. teeing: два агрегата за один проход (Java 12+)
record Summary(long count, BigDecimal total) {}
Summary summary = orders.stream()
        .collect(Collectors.teeing(
                Collectors.counting(),
                Collectors.reducing(BigDecimal.ZERO, Order::amount, BigDecimal::add),
                Summary::new));

Разница между .filter(p).collect(groupingBy(f)) и .collect(groupingBy(f, filtering(p, toList()))) неочевидна и важна: в первом случае регионы без подходящих заказов исчезнут из map, во втором — останутся с пустыми списками. Отчёты ломаются именно на этом.

Collectors.toMap: три способа выстрелить себе в ногу

// 1. Дубликат ключа -> IllegalStateException: Duplicate key
Map<String, Order> byCustomer = orders.stream()
        .collect(Collectors.toMap(Order::customer, o -> o));

// 2. null-значение -> NullPointerException, даже если ключ уникален
Map<String, String> phones = users.stream()
        .collect(Collectors.toMap(User::id, User::phone));   // phone может быть null

// 3. Порядок вставки теряется: под капотом HashMap

Правильная форма почти всегда четырёхаргументная:

Map<String, Order> byCustomer = orders.stream()
        .collect(Collectors.toMap(
                Order::customer,
                Function.identity(),
                (older, newer) -> newer,      // явное правило разрешения конфликта
                LinkedHashMap::new));         // предсказуемый порядок

Про то, почему HashMap не гарантирует порядок и как это связано с хешированием, — в коллекциях и дженериках.

Контракт reduce, который нельзя нарушать

identity обязан удовлетворять combiner(identity, u) == u, а accumulator и combiner — быть ассоциативными и без побочных эффектов. Вычитание неассоциативно: reduce(0, (a, b) -> a - b) даст разный ответ на последовательном и параллельном стриме, и это не баг JDK. Конкатенация reduce("", (a, b) -> a + b) ассоциативна, но квадратична — для строк есть Collectors.joining() поверх StringBuilder. Трёхаргументный reduce(identity, accumulator, combiner) нужен, когда тип аккумулятора отличается от типа элемента, но в 90% случаев вместо него читабельнее collect.

Часть 5. Примитивные стримы и цена боксинга

// плохо: миллион объектов Integer, каждый ~16 байт + ссылка
int sum1 = list.stream().map(Order::items).reduce(0, Integer::sum);

// хорошо: ни одной аллокации
int sum2 = list.stream().mapToInt(Order::items).sum();

// статистика за один проход
IntSummaryStatistics stats = orders.stream()
        .mapToInt(Order::items)
        .summaryStatistics();
System.out.printf("n=%d min=%d avg=%.2f max=%d%n",
        stats.getCount(), stats.getMin(), stats.getAverage(), stats.getMax());

IntStream, LongStream, DoubleStream — единственные примитивные стримы в JDK (CharStream/ByteStream/FloatStream нет, их роль играет IntStream: String.chars() возвращает именно его). Переходы: boxed(), mapToObj(), mapToInt(), asLongStream(). Ловушка: IntStream.sum() возвращает int и молча переполняется — для сумм денег и счётчиков берите mapToLong(...).sum().

Часть 6. Параллельные стримы

Как это работает

parallel() — это переключение флага на конвейере. При терминальной операции конвейер рекурсивно делит источник через trySplit, раздаёт куски в ForkJoinPool.commonPool(), каждый воркер прогоняет свой экземпляр цепочки Sink над своим куском, а затем результаты сливаются combiner-ом.

Как параллельный стрим делит источник через trySplit и склеивает результаты

Размер common pool по умолчанию — availableProcessors() - 1 плюс вызывающий поток. Меняется свойством java.util.concurrent.ForkJoinPool.common.parallelism, но менять его глобально почти всегда неправильно: пул общий на всю JVM.

Когда параллелить

Модель NQ Брайана Гоетца: параллелизм окупается, когда произведение числа элементов N на стоимость обработки одного элемента Q превышает примерно 10 000 машинных операций. Это не формула, а порядок величины.

Чек-лист перед parallel():

  1. Источник делится дёшево? Массив, ArrayList, IntStream.range, HashMap — да. LinkedList, Stream.iterate, BufferedReader.lines, Files.lines — нет (у последних trySplit буферизует, что часто убивает весь выигрыш).
  2. Операции без блокировок и без I/O? Блокирующий вызов в common pool парализует общий пул на всю JVM. Для I/O нужны виртуальные потоки — см. конкурентность.
  3. Нет разделяемого мутабельного состояния?
  4. Порядок не важен? Если да — добавьте .unordered(), это заметно ускорит distinct, limit, skip.
  5. Вы измерили? JMH, не System.nanoTime() в main. См. бенчмаркинг и производительность Java.

Как ломаются параллельные стримы

// ГОНКА: ArrayList не потокобезопасен. Результат — потерянные элементы,
// ArrayIndexOutOfBoundsException или null внутри списка. Воспроизводится не всегда.
List<Integer> bad = new ArrayList<>();
IntStream.range(0, 1_000_000).parallel().forEach(bad::add);

// Правильно: пусть контейнер собирает сам стрим
List<Integer> good = IntStream.range(0, 1_000_000).boxed().toList();

// forEach в параллели НЕ сохраняет порядок; forEachOrdered сохраняет,
// но сериализует хвост конвейера и часто убивает выигрыш
list.parallelStream().forEachOrdered(System.out::println);

Если параллельный стрим действительно нужен в изолированном пуле (чтобы не занимать общий), запускайте его как задачу своего ForkJoinPool — конвейер подхватит текущий пул:

try (ForkJoinPool pool = new ForkJoinPool(8)) {          // Java 19+: FJP — AutoCloseable
    long matched = pool.submit(() ->
            data.parallelStream().filter(Heavy::check).count()
    ).get();
}

Ещё одна тихая ловушка: groupingBy в параллели создаёт временные map на каждый чанк и сливает их. Для больших объёмов есть groupingByConcurrent (требует unordered-стрим), который пишет в один ConcurrentHashMap без слияния.

Часть 7. Gatherers — то, чего не хватало десять лет

Промежуточные операции были закрытым списком: свою добавить было нельзя. Нужно «скользящее окно», «группировать подряд идущие», «distinct по ключу» — пишите самодельный Spliterator или выходите из стрима. JEP 485: Stream Gatherers закрыл пробел: превью в JDK 22–23, финальная версия в JDK 24, доступна в LTS-релизе 25.

// скользящее окно
List<List<Integer>> windows = Stream.of(1, 2, 3, 4, 5)
        .gather(Gatherers.windowSliding(2))
        .toList();
// [[1, 2], [2, 3], [3, 4], [4, 5]]

// накопительная сумма (running total)
List<Integer> running = Stream.of(1, 2, 3, 4)
        .gather(Gatherers.scan(() -> 0, Integer::sum))
        .toList();
// [1, 3, 6, 10]

// параллельная обработка с ограничением, поверх виртуальных потоков,
// с СОХРАНЕНИЕМ порядка — то, чего parallel() никогда не умел для I/O
List<Response> responses = urls.stream()
        .gather(Gatherers.mapConcurrent(16, http::fetch))
        .toList();

Свой gatherer пишется почти так же, как коллектор:

/** distinct по произвольному ключу — то, что все годами писали через Set вручную. */
static <T, K> Gatherer<T, Set<K>, T> distinctBy(Function<? super T, K> key) {
    return Gatherer.ofSequential(
            HashSet::new,
            Gatherer.Integrator.of((seen, item, downstream) ->
                    !seen.add(key.apply(item)) || downstream.push(item)));
}

List<Order> unique = orders.stream()
        .gather(distinctBy(Order::customer))
        .toList();

mapConcurrent — самое важное здесь: он делает то, ради чего люди годами неправильно использовали parallel() с блокирующими вызовами. Механику виртуальных потоков разберём в следующей статье.

Часть 8. Честно о границах

Где Java-стримы проигрывают

  • Проверяемые исключения. Функциональные интерфейсы JDK их не объявляют:
List<String> contents = paths.stream()
        .map(p -> Files.readString(p))   // ошибка компиляции: unhandled IOException
        .toList();

Приходится оборачивать в UncheckedIOException, заводить свой «бросающий» функциональный интерфейс или писать обычный цикл. В C# checked-исключений нет вообще, и этой проблемы там не существует — см. основы C#.

  • Нет ленивых генераторов. В C# yield return даёт ленивый источник в три строки; в Java нужен Spliterator на двадцать. Прямое следствие отсутствия корутин на уровне языка.

  • Нет zip. Операции «склеить два стрима попарно» в JDK нет (была в ранних сборках Java 8, убрали). Обходятся IntStream.range(0, n).mapToObj(...) или Guava Streams.zip.

  • Отладка. Стек-трейс из стрима — это десяток кадров ReferencePipeline$3$1.accept и имена вида Foo$$Lambda/0x00007f.... Спасает встроенный в IntelliJ IDEA Stream Trace, показывающий, что произошло на каждом шаге.

  • Многословность. Collectors.groupingBy(f, Collectors.mapping(g, Collectors.toList())) против Enum.group_by в Elixir или |>-конвейера — это неудобно, и притворяться иначе не стоит.

Где стримы не нужны

Не выигрывают ничего и проигрывают в читаемости:

// хуже
IntStream.range(0, arr.length).forEach(i -> arr[i] = arr[i] * 2);
// лучше
for (int i = 0; i < arr.length; i++) arr[i] *= 2;

// хуже
list.stream().forEach(this::process);
// лучше
for (T item : list) process(item);

// хуже: побочный эффект в map — нарушение контракта, ломается в параллели
orders.stream().map(o -> { audit.log(o); return o.id(); }).toList();

Правило: стрим оправдан, когда описывает преобразование данных в результат. Как только конвейер существует ради побочных эффектов — берите цикл.

Производительность против цикла

Обобщённо, по многократно воспроизведённым замерам JMH: на коллекциях в десятки-сотни элементов простой for быстрее стрима в 1.5–3 раза (стоимость сборки конвейера не амортизируется); на десятках тысяч и больше разница уходит в шум — JIT инлайнит цепочку Sink в один цикл, а escape-анализ убирает аллокации; примитивные стримы практически догоняют цикл, а боксящие Stream<Integer> проигрывают устойчиво и заметно.

Два системных эффекта, о которых стоит знать:

  1. Profile pollution / мегаморфизм. Внутренние классы стримов (ReferencePipeline, Sink) общие на всё приложение. Если через одно и то же место прошли десятки разных лямбд, inline cache вырождается в мегаморфный вызов и JIT перестаёт инлайнить: конвейер, быстрый в микробенчмарке, оказывается медленным в проде.
  2. Предел глубины инлайнинга. Очень длинные конвейеры упираются в MaxInlineLevel, и хвост перестаёт сливаться в один цикл.

Оба эффекта означают одно: микробенчмарк стрима не переносится на прод напрямую. Мерить надо приложение целиком — см. производительность и профилирование.

Часть 9. Сквозной пример

Задача: по логу заказов посчитать топ-3 региона по выручке, с числом заказов и средним чеком, за один проход по данным, без загрузки файла в память.

import java.io.IOException;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.nio.file.*;
import java.util.*;
import java.util.stream.*;

public final class RevenueReport {

    record Order(String region, BigDecimal amount) {}
    record RegionStats(long count, BigDecimal total) {
        BigDecimal average() {
            return count == 0 ? BigDecimal.ZERO
                    : total.divide(BigDecimal.valueOf(count), 2, RoundingMode.HALF_UP);
        }
    }

    public static void main(String[] args) throws IOException {
        Path path = Path.of("orders.csv");

        // try-with-resources ОБЯЗАТЕЛЕН: Files.lines держит открытый дескриптор файла
        try (Stream<String> lines = Files.lines(path)) {

            Map<String, RegionStats> byRegion = lines
                    .skip(1)                                        // заголовок CSV
                    .map(RevenueReport::parse)
                    .flatMap(Optional::stream)                      // битые строки просто выпадают
                    .collect(Collectors.groupingBy(
                            Order::region,
                            Collectors.teeing(                      // два агрегата за один проход
                                    Collectors.counting(),
                                    Collectors.reducing(BigDecimal.ZERO, Order::amount, BigDecimal::add),
                                    RegionStats::new)));

            byRegion.entrySet().stream()
                    .sorted(Map.Entry.<String, RegionStats>comparingByValue(
                            Comparator.comparing(RegionStats::total)).reversed())
                    .limit(3)
                    .forEach(e -> System.out.printf("%-10s заказов=%-6d выручка=%-12s чек=%s%n",
                            e.getKey(), e.getValue().count(),
                            e.getValue().total(), e.getValue().average()));
        }
    }

    /** Разбор строки CSV: region,amount. Пустой Optional для битых строк. */
    private static Optional<Order> parse(String line) {
        String[] parts = line.split(",", -1);
        if (parts.length < 2) return Optional.empty();
        try {
            return Optional.of(new Order(parts[0].strip(), new BigDecimal(parts[1].strip())));
        } catch (NumberFormatException e) {
            return Optional.empty();
        }
    }
}
EMEA       заказов=18432  выручка=9184320.55  чек=498.28
APAC       заказов=12007  выручка=6403110.00  чек=533.28
LATAM      заказов=7741   выручка=3120884.10  чек=403.16

Принципиальное здесь: Files.lines читает файл лениво, строка за строкой, и память не зависит от размера файла; Optional::stream заменяет пару «filter + map» и заодно отбрасывает битые строки; teeing считает count и sum за один обход вместо двух; сортировка применяется к результату группировки — к десяткам регионов, а не к миллионам строк; try-with-resources обязателен, иначе дескриптор течёт и вы упрётесь в Too many open files.

Типичные ошибки: чек-лист

  1. Повторное использование стримаIllegalStateException. Храните Supplier<Stream<T>>.
  2. Files.lines / Files.walk без try-with-resources — утечка файловых дескрипторов.
  3. peek для полезной работы — с Java 9 может быть пропущен полностью вместе с count().
  4. Побочные эффекты в map/filter — ломается при parallel(), нарушает контракт API.
  5. Мутация внешней коллекции в forEach вместо collect — гонка в параллели.
  6. Collectors.toMap без merge-функцииIllegalStateException на дубликате ключа; с null-значениемNullPointerException без внятного сообщения.
  7. Stream<Integer> вместо IntStream в горячем коде — боксинг миллионов объектов.
  8. IntStream.sum() для больших сумм — тихое переполнение int; берите mapToLong.
  9. comparing(...).reversed() разворачивает всю цепочку компараторов, а не последний ключ.
  10. sorted() или distinct() на бесконечном стриме — зависание без ошибки.
  11. parallel() с блокирующим I/O — общий ForkJoinPool встаёт на всю JVM.
  12. parallel() на LinkedList, Stream.iterate, BufferedReader.lines — источник не делится.
  13. filter вместо Collectors.filtering внутри groupingBy — исчезают пустые группы.
  14. Захват this долгоживущей лямбдой — объект не собирается GC.
  15. Массив-на-один-элемент как обход effectively final — работает и лжёт про потокобезопасность.
  16. Длинный конвейер вместо цикла там, где нет преобразования данных.
  17. Вера в микробенчмарк — мегаморфизм call site в проде меняет картину.

Мини-итог

  • Лямбда — это invokedynamic плюс LambdaMetafactory, а не анонимный класс. Без захвата экземпляр переиспользуется, с захватом создаётся новый; цена лямбд — не в вызове, а в старте.
  • Стрим не хранит данные, одноразов и ленив: промежуточные операции только строят цепочку стадий, работа начинается с терминальной.
  • Fusion означает, что элемент проходит весь конвейер целиком, а не «слой за слоем». Отсюда и короткое замыкание, и константная память.
  • Stateful-операции (sorted, distinct) — барьеры: отменяют fusion и требуют O(N) памяти. Ставьте их как можно позже в конвейере.
  • Всё поведение источника задаёт Spliterator: делимость и характеристики (SIZED, ORDERED, DISTINCT) прямо включают и выключают оптимизации.
  • Коллектор — это пятёрка supplier/accumulator/combiner/finisher/characteristics. Свой коллектор пишется в десять строк и часто убирает целый проход по данным.
  • parallel() — не «сделать быстрее», а «разделить, посчитать, склеить»: нужны делимый источник, ассоциативный combiner, отсутствие I/O и разделяемого состояния. Без всех четырёх условий он вредит. Gatherers (JDK 24+) сделали промежуточные операции расширяемыми, а mapConcurrent закрыл давнюю дыру с параллельным I/O.
  • Стрим оправдан там, где описывает преобразование данных в результат. Для побочных эффектов и коротких коллекций обычный цикл честнее и быстрее.

Источники

Что дальше

Мы разобрали функциональный стиль в однопоточном мире и на краю заглянули в параллельный: parallel(), ForkJoinPool, mapConcurrent. Но параллельный стрим — узкий частный случай конкурентности. Дальше — полная картина: как устроены потоки платформы, чем executors лучше new Thread(), как композировать асинхронные операции через CompletableFuture и почему виртуальные потоки Project Loom меняют подход к серверному коду радикальнее, чем в своё время это сделали стримы.

Конкурентность: потоки, executors, CompletableFuture, виртуальные потоки

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов