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 залинкован и стоит ровно как обычный вызов.
метафабрика больше не участвует App->>CS: второе и далее CS-->>App: экземпляр (или тот же самый, если захвата нет)
Вывод про старт приложения. Генерация 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; именно оно делает стримы конкурентоспособными с циклом.
Анатомия конвейера
Collection / массив
Files.lines / generate"] --> H["Head-стадия
обёртка над Spliterator"] H --> I1["filter
stateless"] I1 --> I2["map
stateless"] I2 --> I3["sorted
STATEFUL — барьер"] I3 --> I4["limit
short-circuit"] I4 --> T["Терминальная операция
collect / reduce / forEach"] T --> R["Результат или побочный эффект"] subgraph L["Ленивая часть: только строит связный список стадий"] H I1 I2 I3 I4 end style L stroke-dasharray: 5 5
Три вида операций:
- Промежуточные stateless —
map,filter,flatMap,mapMulti,peek,mapToInt. Обрабатывают элемент, не помня о других. Сливаются в один проход. - Промежуточные stateful —
sorted,distinct,limit,skip,takeWhile,dropWhile. Требуют буфера или счётчика.sortedиdistinct— полноценные барьеры: до них конвейер обязан прочитать весь вход. - Терминальные — запускают выполнение. Часть из них короткозамкнутые:
findFirst,findAny,anyMatch,allMatch,noneMatch, а из промежуточных —limitиtakeWhile.
возвращает новую стадию Собран --> Выполняется: терминальная операция Выполняется --> Потреблён: результат получен Собран --> Связан: sourceStage.linkedOrConsumed = true Связан --> Ошибка: повторное использование Потреблён --> Ошибка: повторная терминальная операция Ошибка --> [*]: IllegalStateException
stream has already been operated upon or closed Потреблён --> [*]
Отсюда правило: стрим нельзя сохранить в поле и переиспользовать. Если нужно
«переиспользовать конвейер» — храните 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-ом.
Размер common pool по умолчанию — availableProcessors() - 1 плюс вызывающий поток.
Меняется свойством java.util.concurrent.ForkJoinPool.common.parallelism, но менять его
глобально почти всегда неправильно: пул общий на всю JVM.
Когда параллелить
Модель NQ Брайана Гоетца: параллелизм окупается, когда произведение числа элементов N на стоимость обработки одного элемента Q превышает примерно 10 000 машинных операций. Это не формула, а порядок величины.
Чек-лист перед parallel():
- Источник делится дёшево? Массив,
ArrayList,IntStream.range,HashMap— да.LinkedList,Stream.iterate,BufferedReader.lines,Files.lines— нет (у последнихtrySplitбуферизует, что часто убивает весь выигрыш). - Операции без блокировок и без I/O? Блокирующий вызов в common pool парализует общий пул на всю JVM. Для I/O нужны виртуальные потоки — см. конкурентность.
- Нет разделяемого мутабельного состояния?
- Порядок не важен? Если да — добавьте
.unordered(), это заметно ускоритdistinct,limit,skip. - Вы измерили? 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(...)или GuavaStreams.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> проигрывают устойчиво и заметно.
Два системных эффекта, о которых стоит знать:
- Profile pollution / мегаморфизм. Внутренние классы стримов (
ReferencePipeline,Sink) общие на всё приложение. Если через одно и то же место прошли десятки разных лямбд, inline cache вырождается в мегаморфный вызов и JIT перестаёт инлайнить: конвейер, быстрый в микробенчмарке, оказывается медленным в проде. - Предел глубины инлайнинга. Очень длинные конвейеры упираются в
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.
Типичные ошибки: чек-лист
- Повторное использование стрима →
IllegalStateException. ХранитеSupplier<Stream<T>>. Files.lines/Files.walkбез try-with-resources — утечка файловых дескрипторов.peekдля полезной работы — с Java 9 может быть пропущен полностью вместе сcount().- Побочные эффекты в
map/filter— ломается приparallel(), нарушает контракт API. - Мутация внешней коллекции в
forEachвместоcollect— гонка в параллели. Collectors.toMapбез merge-функции —IllegalStateExceptionна дубликате ключа; с null-значением —NullPointerExceptionбез внятного сообщения.Stream<Integer>вместоIntStreamв горячем коде — боксинг миллионов объектов.IntStream.sum()для больших сумм — тихое переполнениеint; беритеmapToLong.comparing(...).reversed()разворачивает всю цепочку компараторов, а не последний ключ.sorted()илиdistinct()на бесконечном стриме — зависание без ошибки.parallel()с блокирующим I/O — общий ForkJoinPool встаёт на всю JVM.parallel()наLinkedList,Stream.iterate,BufferedReader.lines— источник не делится.filterвместоCollectors.filteringвнутриgroupingBy— исчезают пустые группы.- Захват
thisдолгоживущей лямбдой — объект не собирается GC. - Массив-на-один-элемент как обход
effectively final— работает и лжёт про потокобезопасность. - Длинный конвейер вместо цикла там, где нет преобразования данных.
- Вера в микробенчмарк — мегаморфизм 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.- Стрим оправдан там, где описывает преобразование данных в результат. Для побочных эффектов и коротких коллекций обычный цикл честнее и быстрее.
Источники
- Package java.util.stream (Javadoc) — раздел про ленивость, порядок и параллелизм обязателен к прочтению целиком.
- Javadoc Collector — формальные требования к supplier/accumulator/combiner.
- Javadoc Spliterator — семантика характеристик и
trySplit. - Brian Goetz, Translation of Lambda Expressions — как именно компилируются лямбды.
- Brian Goetz, State of the Lambda: Libraries Edition — почему API спроектирован так, а не иначе.
- Brian Goetz, Java Streams, часть 3: параллелизм — модель NQ и разбор, когда параллель окупается.
- JEP 107: Bulk Data Operations for Collections и JEP 485: Stream Gatherers — начало и последняя глава истории Stream API.
- JLS §15.27: Lambda Expressions — формальные правила, включая захват и
this. - Исходники AbstractPipeline в OpenJDK — лучший способ понять fusion: читайте
wrapSinkиcopyInto. - Urma, Fusco, Mycroft, Modern Java in Action — самая полная книга именно про стримы и лямбды.
- Bloch, Effective Java, 3rd ed. — главы 6 и 7 (пункты 42–48) про лямбды и стримы.
- Doug Lea, A Java Fork/Join Framework — механика пула, на котором работает
parallel(). - JMH — единственный корректный способ померить всё вышесказанное на своей нагрузке.
Что дальше
Мы разобрали функциональный стиль в однопоточном мире и на краю заглянули в параллельный:
parallel(), ForkJoinPool, mapConcurrent. Но параллельный стрим — узкий частный случай
конкурентности. Дальше — полная картина: как устроены потоки платформы, чем executors лучше
new Thread(), как композировать асинхронные операции через CompletableFuture и почему
виртуальные потоки Project Loom меняют подход к серверному коду радикальнее, чем в своё
время это сделали стримы.
Конкурентность: потоки, executors, CompletableFuture, виртуальные потоки