Java for Beginner
870 subscribers
1.01K photos
275 videos
14 files
1.69K links
Канал от новичков для новичков!
Изучайте Java вместе с нами!
Здесь мы обмениваемся опытом и постоянно изучаем что-то новое!

Наш YouTube канал - https://www.youtube.com/@Java_Beginner-Dev

Наш канал на RUTube - https://rutube.ru/channel/37896292/
Download Telegram
flatMap для раскрытия вложенности

Когда Optional содержит объект, который сам возвращает поток или коллекцию, flatMap объединяет уровни вложенности:
// Найти все заказы пользователя по ID, если пользователь существует
Stream<Order> userOrders = userRepository.findById(userId)
.stream() // Stream<User>
.flatMap(u -> u.getOrders().stream()); // Stream<Order>

// Сбор в список, если пользователь найден
List<Order> orders = userOrders.collect(toList());


Здесь flatMap играет двойную роль: раскрывает Optional (через stream()) и раскрывает внутреннюю коллекцию заказов. Пустой Optional прерывает цепочку, возвращая пустой поток.

Более сложный сценарий — цепочка зависимых Optional:
// Было: глубокая вложенность ifPresent
optionalUser.ifPresent(user ->
user.getProfile().ifPresent(profile ->
profile.getAvatar().ifPresent(avatar ->
process(avatar)
)
)
);

// Стало: плоская композиция через stream и flatMap
userRepository.findById(userId)
.stream()
.flatMap(u -> u.getProfile().stream())
.flatMap(p -> p.getAvatar().stream())
.forEach(this::process);


Каждый flatMap обрабатывает потенциальное отсутствие значения, пропуская дальнейшие операции для пустых Optional. Результат — линейная цепочка вместо пирамиды вложенности.


Stream в Optional: reduce и find

Обратное направление — терминальные операции потока, возвращающие Optional. Это естественное завершение конвейера, где результат может отсутствовать:
Optional<Book> oldest = library.stream()
.filter(b -> b.genre() == Genre.FICTION)
.min(comparing(Book::year)); // Optional<Book>

// Интеграция обратно в поток
List<String> oldestTitle = oldest
.stream() // Optional -> Stream
.map(Book::title)
.collect(toList());


findFirst(), findAny(), min(), max(), reduce() — все возвращают Optional, сигнализируя о возможности пустого результата. Композиция этих операций с Optional.stream() создаёт бесшовные цепочки обработки, где отсутствие данных — нормальный сценарий, не исключительная ситуация.


Обработка коллекций в Optional

Когда Optional содержит коллекцию, которая сама может быть пустой, Optional.stream() сочетается с flatMap для нормализации:
// Репозиторий возвращает Optional<List<Order>> — неудачный дизайн, но встречается
Optional<List<Order>> optionalOrders = orderRepository.findByCustomerId(id);

// Нормализация: Optional<List<Order>> -> Stream<Order>
Stream<Order> orders = optionalOrders
.stream() // Stream<List<Order>>
.flatMap(List::stream); // Stream<Order>, пустой если Optional пуст или список пуст

// Или в одном выражении
List<Order> result = optionalOrders.stream()
.flatMap(List::stream)
.filter(Order::isActive)
.collect(toList());


Этот паттерн устраняет необходимость проверять optionalOrders.isPresent() и optionalOrders.get().isEmpty() отдельно.


#Java #для_новичков #beginner #stream_api #optional
👍4
Раздел 8. Stream API и функциональный стиль в Java

Глава 7: Stream API в экосистеме Java. Интеграционные паттерны

Рефакторинг legacy-кода в проекте «Библиотека»

Представьте, что в проекте «Библиотека» накопился такой метод — подсчёт статистики по авторам с множеством условий:
public Map<String, AuthorStats> calculateAuthorStatsLegacy(List<Book> books, int minYear) {
Map<String, AuthorStats> result = new HashMap<>();

for (Book book : books) {
// Пропускаем старые книги
if (book.getYear() < minYear) {
continue;
}

// Пропускаем без жанров
if (book.getGenres().isEmpty()) {
continue;
}

// Пропускаем, если нет цены
Double price = book.getPrice();
if (price == null || price <= 0) {
continue;
}

String author = book.getAuthor();
AuthorStats stats = result.get(author);

if (stats == null) {
stats = new AuthorStats();
stats.firstBookYear = book.getYear();
result.put(author, stats);
}

// Обновляем статистику
stats.bookCount++;
stats.totalPages += book.getPages();
stats.totalPrice += price;

// Обновляем самую старую книгу
if (book.getYear() < stats.firstBookYear) {
stats.firstBookYear = book.getYear();
}

// Собираем уникальные жанры
for (String genre : book.getGenres()) {
if (!stats.genres.contains(genre)) {
stats.genres.add(genre);
}
}
}

// Пост-обработка: считаем среднее
for (AuthorStats stats : result.values()) {
stats.averagePrice = stats.totalPrice / stats.bookCount;
}

return result;
}

static class AuthorStats {
int bookCount;
int totalPages;
double totalPrice;
double averagePrice;
int firstBookYear;
Set<String> genres = new HashSet<>();
}


Характеристики кода:
40+ строк в одном методе
4 уровня вложенности в пике
Смешение фильтрации, трансформации и агрегации
Два прохода: основной цикл + пост-обработка


#Java #для_новичков #beginner #stream_api #практика
👍4
Мутация состояния на каждом шагу

Шаг 1: Выделить источник данных в stream

Не меняем логику — только оборачиваем в stream:
public Map<String, AuthorStats> calculateAuthorStatsStep1(List<Book> books, int minYear) {
Map<String, AuthorStats> result = new HashMap<>();

books.stream().forEach(book -> { // Заменили for-each на forEach
// ... весь прежний код без изменений ...
});

// пост-обработка
result.values().forEach(stats ->
stats.averagePrice = stats.totalPrice / stats.bookCount
);

return result;
}


Что достигли: Начали привыкать к синтаксису. Никаких функциональных преимуществ ещё нет — это подготовка.


Шаг 2: Заменить простые фильтры

Выносим continue как предикаты filter. Каждый фильтр — отдельная проверка:
public Map<String, AuthorStats> calculateAuthorStatsStep2(List<Book> books, int minYear) {
Map<String, AuthorStats> result = new HashMap<>();

books.stream()
.filter(book -> book.getYear() >= minYear) // было: if (year < minYear) continue
.filter(book -> !book.getGenres().isEmpty()) // было: if (genres.isEmpty()) continue
.filter(book -> { // было: if (price == null || price <= 0) continue
Double price = book.getPrice();
return price != null && price > 0;
})
.forEach(book -> {
// ... урезанный цикл без фильтров ...
String author = book.getAuthor();
AuthorStats stats = result.computeIfAbsent(author, k -> {
AuthorStats s = new AuthorStats();
s.firstBookYear = book.getYear();
return s;
});

stats.bookCount++;
stats.totalPages += book.getPages();
stats.totalPrice += book.getPrice();

if (book.getYear() < stats.firstBookYear) {
stats.firstBookYear = book.getYear();
}

book.getGenres().forEach(genre -> stats.genres.add(genre));
});

result.values().forEach(stats ->
stats.averagePrice = stats.totalPrice / stats.bookCount
);

return result;
}


Улучшения:
Фильтры вынесены, читаются как декларативные условия
Заменили if (stats == null) на computeIfAbsent
Жанры теперь через forEach
Проблема осталась: Сложная логика агрегации всё ещё внутри forEach с мутацией.


#Java #для_новичков #beginner #stream_api #практика
👍4
Шаг 3: Вынести преобразование в map

Выделяем «чистые» данные для агрегации — промежуточный объект:
// Вспомогательный record для промежуточных данных
record BookData(String author, int year, int pages, double price, List<String> genres) {}

public Map<String, AuthorStats> calculateAuthorStatsStep3(List<Book> books, int minYear) {
Map<String, List<BookData>> grouped = books.stream()
.filter(book -> book.getYear() >= minYear)
.filter(book -> !book.getGenres().isEmpty())
.filter(book -> book.getPrice() != null && book.getPrice() > 0)
.map(book -> new BookData(
book.getAuthor(),
book.getYear(),
book.getPages(),
book.getPrice(),
book.getGenres()
))
.collect(Collectors.groupingBy(BookData::author));

// Отдельный проход для агрегации
Map<String, AuthorStats> result = new HashMap<>();
grouped.forEach((author, dataList) -> {
AuthorStats stats = new AuthorStats();
stats.bookCount = dataList.size();
stats.totalPages = dataList.stream().mapToInt(BookData::pages).sum();
stats.totalPrice = dataList.stream().mapToDouble(BookData::price).sum();
stats.averagePrice = stats.totalPrice / stats.bookCount;
stats.firstBookYear = dataList.stream()
.mapToInt(BookData::year)
.min()
.orElse(0);
stats.genres = dataList.stream()
.flatMap(d -> d.genres().stream())
.collect(Collectors.toSet());

result.put(author, stats);
});

return result;
}


Улучшения:
Разделили фильтрацию, группировку и агрегацию
Каждый этап — отдельный, тестируемый блок
Убрали мутацию внутри потока
Новая проблема: Два прохода по данным, избыточные промежуточные структуры.


Шаг 4: Сложная агрегация через collect

Объединяем всё в один проход с кастомным коллектором:
public Map<String, AuthorStats> calculateAuthorStatsStep4(List<Book> books, int minYear) {
return books.stream()
.filter(book -> book.getYear() >= minYear)
.filter(book -> !book.getGenres().isEmpty())
.filter(book -> book.getPrice() != null && book.getPrice() > 0)
.collect(Collectors.groupingBy(
Book::getAuthor,
Collector.of(
AuthorStats::new, // supplier
(stats, book) -> { // accumulator
if (stats.bookCount == 0) {
stats.firstBookYear = book.getYear();
} else {
stats.firstBookYear = Math.min(stats.firstBookYear, book.getYear());
}
stats.bookCount++;
stats.totalPages += book.getPages();
stats.totalPrice += book.getPrice();
stats.genres.addAll(book.getGenres());
},
(left, right) -> { // combiner для parallelStream
if (left.bookCount == 0) return right;
if (right.bookCount == 0) return left;

AuthorStats merged = new AuthorStats();
merged.bookCount = left.bookCount + right.bookCount;
merged.totalPages = left.totalPages + right.totalPages;
merged.totalPrice = left.totalPrice + right.totalPrice;
merged.firstBookYear = Math.min(left.firstBookYear, right.firstBookYear);
merged.genres.addAll(left.genres);
merged.genres.addAll(right.genres);
return merged;
},
stats -> { // finisher
stats.averagePrice = stats.totalPrice / stats.bookCount;
return stats;
}
)
));
}


Достигнуто:
Один проход по данным
Нет промежуточных структур
Потенциально parallelizable
Цена: 30 строк сложного кода вместо 40 строк простого. Стало ли читаемее?


#Java #для_новичков #beginner #stream_api #практика
👍4
Шаг 5: Точка остановки — прагматичный компромисс

Вернёмся к шагу 3 и доработаем его до читаемого состояния без избыточной сложности:
public Map<String, AuthorStats> calculateAuthorStatsPragmatic(List<Book> books, int minYear) {
// Предварительная фильтрация — ясная и тестируемая
List<Book> validBooks = books.stream()
.filter(this::isValidForStats)
.filter(book -> book.getYear() >= minYear)
.collect(Collectors.toList());

// Группировка — стандартная операция
Map<String, List<Book>> byAuthor = validBooks.stream()
.collect(Collectors.groupingBy(Book::getAuthor));

// Агрегация — явный цикл с понятной логикой
Map<String, AuthorStats> result = new HashMap<>();
for (Map.Entry<String, List<Book>> entry : byAuthor.entrySet()) {
result.put(entry.getKey(), aggregateStats(entry.getValue()));
}

return result;
}

private boolean isValidForStats(Book book) {
return !book.getGenres().isEmpty()
&& book.getPrice() != null
&& book.getPrice() > 0;
}

private AuthorStats aggregateStats(List<Book> authorBooks) {
AuthorStats stats = new AuthorStats();
stats.bookCount = authorBooks.size();
stats.totalPages = authorBooks.stream()
.mapToInt(Book::getPages)
.sum();
stats.totalPrice = authorBooks.stream()
.mapToDouble(Book::getPrice)
.sum();
stats.averagePrice = stats.totalPrice / stats.bookCount;
stats.firstBookYear = authorBooks.stream()
.mapToInt(Book::getYear)
.min()
.orElse(0);
stats.genres = authorBooks.stream()
.flatMap(b -> b.getGenres().stream())
.collect(Collectors.toSet());
return stats;
}



Почему это лучше шага 4:

Методы isValidForStats и aggregateStats можно тестировать отдельно
Нет 30-строчного лямбда-ада в коллекторе
Junior-разработчик поймёт код за 5 минут
Легко добавить логирование или метрики в промежуточные точки


Критерии остановки рефакторинга

Остановитесь, если:
Код стал медленнее без причины
Появились неочевидные побочные эффекты
// Опасно: параллельный stream с непотокобезопасным accumulators
.collect(Collectors.groupingByConcurrent(
Book::getAuthor,
Collector.of(
() -> new AuthorStats(), // ОК — новый для каждого
(stats, book) -> stats.genres.addAll(...), // Опасно! HashSet не thread-safe
...
)
))


Теряется отладочная информация
В исходном коде можно было поставить breakpoint на любой continue. В сложном коллекторе — только внутри лямбд.

Команда не владеет паттернами
Если 3 из 5 разработчиков не знают Collector.of — код неподдерживаем.


Практические задания

Задача 1: рефакторинг с остановкой

В проекте «Библиотека» найдите метод с циклом, содержащим:
2+ условия continue
мутацию аккумулятора
вложенные циклы
Проведите рефакторинг до шага 3 (группировка + явная агрегация). Остановитесь. Обоснуйте выбор.

Задача 2: сравнение читаемости

Покажите шаг 4 (полный коллектор) коллеге, не знакомому с Stream API. Засеките время, за которое он поймёт логику. Повторите с шагом 5. Зафиксируйте разницу.

Задача 3: добавление функциональности

К обоим вариантам (шаг 4 и шаг 5) добавьте требование: «пропускать авторов с менее чем 3 книгами».
В каком варианте изменение проще? В каком меньше риск регрессии?

Задача 4: документация компромисса (звёздочка)
Создайте в проекте docs/stream-guidelines.md с правилами:
когда использовать Stream
когда остановиться
примеры «хорошего», «плохого» и «достаточного» кода из вашей кодовой базы


#Java #для_новичков #beginner #stream_api #практика
👍4
Раздел 8. Stream API и функциональный стиль в Java

Глава 8: За пределами коллекций. Бесконечность и I/O

Бесконечные стримы и ленивая генерация

Ранее Stream API рассматривался как инструмент обработки конечных коллекций данных. Но фундаментальная мощь потоков — в их способности моделировать бесконечные последовательности, вычисляемые по требованию. Это сдвигает фокус с "данных в памяти" к "процессу генерации", открывая паттерны для числовых последовательностей, генераторов уникальных идентификаторов, обработки потоковых источников вроде сетевых соединений.


Stream.iterate: рекурсия с состоянием

Метод Stream.iterate создаёт поток, где каждый следующий элемент вычисляется из предыдущего.

Перегрузка Java 8 принимает начальное значение и унарную функцию:
// Натуральные числа: 0, 1, 2, 3...
Stream<Integer> natural = Stream.iterate(0, n -> n + 1);

// Степени двойки: 1, 2, 4, 8...
Stream<Integer> powersOfTwo = Stream.iterate(1, n -> n * 2);

// Фибоначчи: пара (a, b) -> (b, a+b)
Stream<long[]> fibonacci = Stream.iterate(
new long[]{0, 1},
f -> new long[]{f[1], f[0] + f[1]}
);


Ключевая особенность: iterate сохраняет состояние между элементами. Это делает его непараллелизуемым — каждый элемент зависит от предыдущего, разделение невозможно. Попытка вызвать .parallel() на таком потоке не даст выигрыша, а может замедлить из-за накладных расходов.


Java 9+: iterate с предикатом остановки

Перегрузка с тремя параметрами добавляет условие продолжения:
// Числа от 0 до 99
Stream<Integer> limited = Stream.iterate(
0, // seed
n -> n < 100, // hasNext (предикат продолжения)
n -> n + 1 // next (функция следующего)
);


Это заменяет паттерн iterate(...).limit(n) более семантически ясным конструктом. Предикат проверяется перед генерацией каждого элемента, включая seed. Если hasNext(seed) ложно, поток пуст.


Stream.generate: чистая генерация без состояния

Stream.generate принимает Supplier — функцию без аргументов, возвращающую значение. Каждый вызов независим, что теоретически позволяет параллелизм (хотя реализация в OpenJDK не делит такие потоки эффективно).
// Постоянное значение
Stream<String> constants = Stream.generate(() -> "repeat");

// Случайные числа
Stream<Double> randoms = Stream.generate(Math::random);

// Уникальные ID через атомарный счётчик
AtomicLong idGenerator = new AtomicLong(0);
Stream<String> uniqueIds = Stream.generate(() -> "ID-" + idGenerator.incrementAndGet());


generate идеален для стохастических или внешне определяемых последовательностей, где нет рекуррентной зависимости. Но бесконечность требует ограничения: limit, takeWhile или короткое замыкание терминальной операцией.


takeWhile и dropWhile: предикатные границы

Java 9 добавила методы takeWhile и dropWhile, критически отличные от filter по семантике:
filter проверяет каждый элемент независимо, пропуская или отбрасывая по условию
takeWhile пропускает элементы, пока предикат истинен, и останавливает поток при первом ложном
dropWhile отбрасывает элементы, пока предикат истинен, и продолжает с первого ложного

// Сортированный поток температур
IntStream temperatures = IntStream.of(15, 18, 22, 25, 19, 12, 8);

// takeWhile: температура ниже 25
temperatures.takeWhile(t -> t < 25)
.forEach(System.out::println); // 15, 18, 22 — остановка на 25

// dropWhile: пропускаем прохладную погоду
temperatures.dropWhile(t -> t < 20)
.forEach(System.out::println); // 22, 25, 19, 12, 8


Ключевое отличие takeWhile от filter — short-circuit поведение. filter(t -> t < 25) обработал бы все элементы, включая 19, 12, 8. takeWhile останавливается, предполагая, что последовательность упорядочена и дальнейшие элементы не интересны.

Это критично для бесконечных потоков:
// Бесконечная последовательность, ограниченная условием
Stream.iterate(1, n -> n * 2)
.takeWhile(n -> n < 1000) // 1, 2, 4, 8, ..., 512 — остановка
.forEach(System.out::println);


Без takeWhile пришлось бы использовать limit, требующий знания количества элементов заранее, или filter, не останавливающийся.


#Java #для_новичков #beginner #stream_api
👍6
Сценарий: генерация уникальных идентификаторов

Комбинация generate с takeWhile или limit создаёт контролируемые генераторы:

// UUID с проверкой уникальности в пределах сессии
Set<String> usedIds = ConcurrentHashMap.newKeySet();
Stream<String> uniqueIds = Stream.generate(UUID::randomUUID)
.map(UUID::toString)
.filter(usedIds::add) // add возвращает false если уже есть, фильтруем дубликаты
.limit(1000); // строго 1000 уникальных

// Или с takeWhile для условной остановки
Stream<String> idsUntilPattern = Stream.generate(this::generateId)
.takeWhile(id -> !id.contains("STOP")); // Генерация до спецпаттерна


Атомарные структуры (AtomicLong, ConcurrentHashMap) обеспечивают потокобезопасность при параллельном доступе, хотя generate в стандартной реализации не эффективно параллелится.


Сценарий: чтение потоковых данных

Бесконечные стримы моделируют внешние источники с неизвестным объёмом:
```
// Симуляция чтения из сокета: байты до терминатора
Stream<Byte> socketStream = Stream.generate(() -> readFromSocket())
.takeWhile(b -> b != -1); // -1 как EOF

// Обработка пакетов до специального маркера
Stream<Packet> packetStream = Stream.generate(this::readNextPacket)
.takeWhile(p -> !p.isTerminator());

Важно: такие потоки требуют управления ресурсами. Stream.generate не знает о необходимости закрыть сокет. Использование в try-with-resources или явное закрытие источника в Supplier — ответственность разработчика.


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

Сравним подходы к генерации последовательностей:
// Плохо: немедленное создание списка
List<Integer> eager = new ArrayList<>();
for (int i = 0; i < 1000000; i++) {
if (isPrime(i)) eager.add(i);
}

// Хорошо: ленивый поток, элементы вычисляются по требованию
IntStream primes = IntStream.iterate(2, n -> n + 1)
.filter(this::isPrime)
.takeWhile(n -> n < 1000000);

// Потребление только первых 10, остальные никогда не вычислены
primes.limit(10).forEach(System.out::println);


Ленивость позволяет работать с "потенциально бесконечными" последовательностями, фактически обрабатывая только необходимый минимум.


Ограничения бесконечных потоков

Параллелизм: iterate не параллелится, generate — плохо. Бесконечные потоки — последовательная абстракция.
Состояние: iterate требует небольшого состояния (предыдущий элемент), но не масштабируется. Сложное состояние в Supplier требует синхронизации.
Ресурсы: бесконечность — концептуальная. Реальные источники (сокеты, файлы) требуют закрытия. Stream API не управляет жизненным циклом внешних ресурсов автоматически.


#Java #для_новичков #beginner #stream_api
👍5
Раздел 8. Stream API и функциональный стиль в Java

Глава 8: За пределами коллекций. Бесконечность и I/O

I/O операции как потоки данных

Stream API расширяет свою абстракцию от коллекций в памяти к внешним источникам данных. Файлы, сетевые соединения, процессы — всё это может быть представлено как Stream<String> или Stream<Byte>, с ленивой подкачкой данных по мере необходимости. Этот подход критичен для обработки данных, не помещающихся в RAM: гигабайтные логи, бесконечные потоки событий, построчная обработка без полной загрузки.


Files.lines: файловый поток строк

Метод Files.lines(Path path) — фабрика потоков для текстовых файлов. Он возвращает Stream<String>, где каждый элемент — строка файла, декодированная в указанной или платформенной кодировке (по умолчанию UTF-8).
// Базовое использование с автоматическим закрытием
try (Stream<String> lines = Files.lines(Path.of("access.log"))) {
long errorCount = lines
.filter(line -> line.contains("ERROR"))
.count();
}


Критически важно: Files.lines возвращает поток, реализующий AutoCloseable. Файловый дескриптор открывается при создании потока и должен быть явно закрыт. Без try-with-resources или явного close() ресурс утекает до сборки мусора, что при высокой частоте операций приводит к исчерпанию дескрипторов ОС (Too many open files).

Паттерн обработки больших файлов:
// Фильтрация, трансформация, запись результата — всё в потоковом режиме
try (Stream<String> lines = Files.lines(Path.of("huge_input.txt"));
BufferedWriter writer = Files.newBufferedWriter(Path.of("filtered_output.txt"))) {

lines.parallel() // Осторожно: см. предупреждение ниже
.filter(line -> line.length() > 100)
.map(String::toUpperCase)
.map(line -> line + System.lineSeparator())
.forEachOrdered(line -> {
try {
writer.write(line);
} catch (IOException e) {
throw new UncheckedIOException(e);
}
});
}


Здесь ни одна строка не хранится в памяти целиком. Files.lines использует BufferedReader с буфером 8192 байт, читая файл блоками и выдавая строки по границам \n или \r\n. Обработка происходит пакетно, с постоянным потреблением памяти независимо от размера файла.


BufferedReader.lines: низкоуровневый контроль

Для специализированных сценариев — настройка размера буфера, обработка кодировок, работа с существующим Reader:
// Кастомная кодировка и буфер
try (BufferedReader reader = new BufferedReader(
new InputStreamReader(
new FileInputStream("legacy.txt"),
StandardCharsets.ISO_8859_1
),
16384 // Удвоенный буфер для последовательного чтения
);
Stream<String> lines = reader.lines()) {

lines.map(this::parseLegacyRecord)
.filter(Objects::nonNull)
.forEach(this::processRecord);
}


BufferedReader.lines() возвращает поток с теми же семантиками ленивости, но без привязки к Path — работает с любым Reader, включая StringReader, CharArrayReader, сетевые потоки через InputStreamReader.


#Java #для_новичков #beginner #stream_api
👍5
Управление ресурсами и исключения

I/O потоки добавляют сложность обработки checked исключений. Лямбды в Stream API не объявляют throws, требуя обёртки:
// Проблема: readLine бросает IOException
Stream<String> lines = Files.lines(path);
lines.map(line -> {
// Ошибка компиляции: unreported exception IOException
return expensiveParser.parse(line);
});


Решения:
Обёртка в runtime exception:
.lines.map(line -> {
try {
return expensiveParser.parse(line);
} catch (IOException e) {
throw new UncheckedIOException(e);
}
})


Извлечение в метод с обёрткой:
private ParsedRecord safeParse(String line) {
try {
return expensiveParser.parse(line);
} catch (IOException e) {
throw new UncheckedIOException(e);
}
}

// В потоке
.lines.map(this::safeParse)


Специализированный коллектор для ошибок (см. Урок 3.3): разделение успешных результатов и ошибок без прерывания потока.


Параллелизм и I/O: катастрофическая комбинация

Предупреждение из предыдущих глав приобретает критическую важность для I/O. parallelStream() на Files.lines или BufferedReader.lines() — антипаттерн с тяжёлыми последствиями:
// КАТАСТРОФА: параллельное чтение файла через common pool
try (Stream<String> lines = Files.lines(Path.of("access.log"))) {
lines.parallel() // Нет выигрыша, есть риск
.map(this::blockingDatabaseLookup) // Блокировка воркера
.collect(toList());
}


Проблемы:
Нет разделения: BufferedReader не поддерживает trySplit(). lines().parallel() не делит файл на сегменты — он создаёт иллюзию параллелизма, фактически синхронизируя доступ к общему Reader.

Блокировка common pool: если map содержит блокирующие операции (запросы к БД, HTTP вызовы, Thread.sleep), все воркеры ForkJoinPool.commonPool() замораживаются. Другие компоненты приложения (CompletableFuture, другие parallelStream) парализуются.

Нарушение порядка: forEach в параллельном потоке не сохраняет порядок строк файла. forEachOrdered требует синхронизации, сводя на нет выигрыш.

Исключение: если обработка строк CPU-bound и дорога (сложный парсинг, криптографические операции), а чтение — отдельная стадия, можно разделить:
// Чтение последовательное, обработка параллельная — но с осторожностью
List<String> batch = new ArrayList<>(BATCH_SIZE);
try (Stream<String> lines = Files.lines(Path.of("huge.txt"))) {
lines.forEach(line -> {
batch.add(line);
if (batch.size() >= BATCH_SIZE) {
processBatchParallel(new ArrayList<>(batch)); // Копия для безопасности
batch.clear();
}
});
if (!batch.isEmpty()) processBatchParallel(batch);
}

private void processBatchParallel(List<String> batch) {
batch.parallelStream()
.map(this::expensiveCpuBoundTransform)
.collect(toList()); // Результат куда-то сохраняется
}


Но лучший подход для I/O-bound задач — CompletableFuture с кастомным пулом или virtual threads (Java 21+):
// Правильно: изолированный пул для блокирующих операций
ExecutorService ioPool = Executors.newFixedThreadPool(50);

try (Stream<String> lines = Files.lines(Path.of("urls.txt"))) {
List<CompletableFuture<Response>> futures = lines
.map(url -> CompletableFuture.supplyAsync(
() -> fetchHttp(url),
ioPool
))
.collect(toList());

List<Response> results = futures.stream()
.map(CompletableFuture::join)
.collect(toList());
} finally {
ioPool.shutdown();
}



#Java #для_новичков #beginner #stream_api
👍4
Паттерн: конвейер ETL без промежуточных файлов

Классический паттерн Extract-Transform-Load реализуется через композицию потоков:
// Извлечение из CSV
try (Stream<String> lines = Files.lines(Path.of("input.csv"));
// Загрузка в выходной файл
BufferedWriter writer = Files.newBufferedWriter(Path.of("output.json"))) {

lines.skip(1) // Пропуск заголовка
.map(this::parseCsvLine) // String -> Record
.filter(Objects::nonNull) // Удаление malformed
.map(this::transformToJson) // Record -> JSON string
.forEach(json -> {
try {
writer.write(json);
writer.newLine();
} catch (IOException e) {
throw new UncheckedIOException(e);
}
});
}


Память потребляется постоянно (размер буфера чтения + одна строка обработки), независимо от размера входного файла. Скорость ограничена I/O диска или сети, а не CPU.


Закрытие и обработка ошибок

При исключении в промежуточной операции поток прерывается, но ресурс в try-with-resources закрывается корректно:
try (Stream<String> lines = Files.lines(Path.of("corrupt.txt"))) {
lines.map(this::parse)
.filter(Objects::nonNull)
.forEach(this::process);
// Если parse бросает RuntimeException на 1000-й строке,
// lines.close() вызывается автоматически
}



Но если исключение происходит в терминальной операции, а промежуточные содержат ресурсы (например, map открывает соединения), требуется явное управление:
// Антипаттерн: ресурс внутри map
lines.map(line -> {
Connection conn = pool.borrow(); // Открытие здесь
return query(conn, line); // Если исключение, conn не возвращается
})

// Правильно: try-with-resources внутри лямбды, или вне потока


#Java #для_новичков #beginner #stream_api
👍4