JavaRush /Курсы /Kotlin SELF /Операторы Flow: map, filter, onEach, catch и backpressure...

Операторы Flow: map, filter, onEach, catch и backpressure

Kotlin SELF
54 уровень , 3 лекция
Открыта

1. Зачем нужны Flow-операторы

Когда видишь Flow впервые, рука тянется сделать так: “ну я же умею for и if — значит, просто напишу всё внутри collect { ... }”. Это работает, но обычно превращается в один огромный обработчик, где перемешаны фильтрация, преобразование, отладочная печать и обработка ошибок. Через неделю этот код выглядит как “я не помню, что хотел сказать этим if”.

Операторы Flow помогают держать голову в порядке: каждый шаг конвейера отвечает за одну простую идею. Получается почти как на заводе: сначала детали приходят, потом проходят через фильтр качества, потом их маркируют, потом кладут в коробку, и только в конце мы получаем готовый продукт.

Схема, к которой мы сегодня будем стремиться, выглядит так:

flowchart LR
    A["Источник flow { emit(...) }"] --> B["filter {...}"]
    B --> C["map {...}"]
    C --> D["onEach {...}"]
    D --> E["catch {...}"]
    E --> F["collect {...}"]

2. Пример: поток событий мониторинга

Чтобы примеры были связаны между собой, сделаем маленькую “консольную систему мониторинга”: у нас есть события (например, "INFO", "WARN", "ERROR"), мы хотим пропускать только важные, преобразовывать их в понятный текст, иногда печатать отладку и не падать всем приложением при одной ошибке парсинга.

Начнём с модели события:

data class MonitorEvent(
    val level: String,
    val message: String
)

И ещё договоримся о формате “сырой строки”, как будто она пришла из лога: "INFO|Service started" или "ERROR|Disk is full".

3. Базовый Flow: источник и collect

Перед операторами полезно почувствовать “голую трубу”: есть источник, который эмитит значения, и есть сборщик. Flow запускается только при collect (то есть поток по сути ленивый / cold).

Вот минимальный поток “сырого текста”:

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    val lines = flow {
        emit("INFO|Service started")
        emit("WARN|Cache miss")
    }

    lines.collect { line ->
        println("raw: $line") // raw: INFO|Service started ...
    }
}

Пока это не “анализатор”, а просто доставка значений. Сейчас мы начнём превращать поток в конвейер.

4. Преобразования: map, filter, onEach

map: преобразование и смена типа

map в Flow по смыслу очень похож на map у коллекций: берём каждый элемент и превращаем его во что-то новое. С Flow идея та же, только элементы приходят “во времени”, и преобразование может быть suspend-дружелюбным.

Сделаем парсер строки в MonitorEvent. Для простоты без Regex, обычный split.

fun parseLine(line: String): MonitorEvent {
    val parts = line.split("|", limit = 2)
    val level = parts[0]
    val message = parts.getOrElse(1) { "" }
    return MonitorEvent(level = level, message = message)
}

Теперь применим map:

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    flow {
        emit("INFO|Service started")
        emit("ERROR|Disk is full")
    }
        .map { parseLine(it) }
        .collect { e ->
            println("${e.level}: ${e.message}")
            // INFO: Service started
            // ERROR: Disk is full
        }
}

Обратите внимание на важную вещь: тип потока поменялся. Был Flow<String>, стал Flow<MonitorEvent>. Это суперсила map: он может менять тип, и это делает конвейер намного выразительнее.

Небольшой лайфхак для чтения кода: если вы потерялись, “какого типа поток сейчас”, вспомните, что map меняет тип, а filter — обычно нет.

filter: оставляем только нужное

Теперь представьте, что событий много, а нас интересуют только "WARN" и "ERROR". Тут включается filter: отбор элементов по предикату (лямбда, которая возвращает true/false). Для Flow смысл ровно тот же: значение либо проходит дальше по конвейеру, либо отбрасывается.

Сделаем фильтрацию после парсинга:

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.filter
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    flow {
        emit("INFO|Service started")
        emit("WARN|Cache miss")
        emit("ERROR|Disk is full")
    }
        .map { parseLine(it) }
        .filter { e -> e.level != "INFO" }
        .collect { e ->
            println("IMPORTANT: ${e.level} ${e.message}")
            // IMPORTANT: WARN Cache miss
            // IMPORTANT: ERROR Disk is full
        }
}

Здесь filter не меняет тип: как был поток MonitorEvent, так и остался. Он только “сужает” набор значений.

onEach: побочные эффекты без изменения данных

Когда вы отлаживаете поток, хочется “подсмотреть”, что происходит на каждом шаге. Первая мысль — засунуть println в map. Но это плохой стиль: map по смыслу должен преобразовывать данные, а печать — это побочный эффект.

onEach как раз для этого: выполнить действие для каждого элемента и пропустить элемент дальше без изменений. Это похоже на forEach, только в середине flow-конвейера.

Мы добавим диагностику: покажем, что именно прошло фильтр.

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.filter
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    flow {
        emit("INFO|Service started")
        emit("WARN|Cache miss")
        emit("ERROR|Disk is full")
    }
        .map { parseLine(it) }
        .filter { it.level != "INFO" }
        .onEach { println("debug: passed filter -> $it") }
        .collect { e ->
            println("ALERT: ${e.message}")
        }
}

Плюс onEach в том, что его легко удалить, когда отладка не нужна, и он не “ломает” смысл map.

Кстати, если вы когда-нибудь увидите println внутри map, это не смертный грех, но это почти всегда сигнал: “код пора чуть-чуть привести в чувство”.

5. Ошибки: catch и границы его действия

catch: обработка ошибок в конвейере

Обработка ошибок — это то, что отличает “пример из учебника” от “программы, которая живёт дольше 15 секунд”. В Kotlin мы уже умеем try/catch и понимаем, что исключение прерывает обычный поток выполнения. В Flow ошибки тоже бывают, но часто удобнее ловить их прямо в конвейере оператором catch.

Сделаем парсер строже: если строка не содержит "|", мы считаем это ошибкой формата.

fun parseLineStrict(line: String): MonitorEvent {
    val parts = line.split("|", limit = 2)
    if (parts.size < 2) error("Bad format: '$line'")
    return MonitorEvent(level = parts[0], message = parts[1])
}

Теперь применим catch и “подменим” ошибку на специальное событие:

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    flow {
        emit("INFO|Service started")
        emit("BROKEN LINE")
        emit("ERROR|Disk is full")
    }
        .map { parseLineStrict(it) }
        .catch { e ->
            println("handled error: ${e.message}") // handled error: Bad format: 'BROKEN LINE'
            emit(MonitorEvent("ERROR", "Parser failed, but app continues"))
        }
        .collect { e ->
            println("event: ${e.level} ${e.message}")
        }
}

Здесь важно понять механику: catch перехватывает ошибки, которые произошли до него в конвейере (в upstream). То есть если ошибка вылетела в map { parseLineStrict(it) }, наш catch её увидит. А если ошибка случится позже (например, внутри collect), то catch может не помочь.

Почему catch не ловит ошибку внутри collect

Это типичная ловушка: вы ставите catch, успокаиваетесь, и думаете, что поток теперь “неубиваемый”. Но catch — не волшебный зонтик от всех бед, он ловит только то, что случилось “выше по трубе”.

Сделаем искусственную ошибку внутри collect:

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    flow {
        emit(1)
    }
        .catch { println("catch: ${it.message}") }
        .collect { _ ->
            error("boom in collect")
        }
}

Этот код упадёт (если не обернуть collect в try/catch), потому что ошибка случилась уже в потребителе. Это логично: collect — терминальная операция, и если именно потребитель падает, то спасать его должен внешний try/catch, как и в обычной жизни.

Практическое правило: catch — для ошибок источника и операторов до collect. Ошибки обработчика в collect — ловим снаружи обычным try/catch.

6. Backpressure в Flow

Слово backpressure звучит как название рок-группы (“Backpressure и их новый альбом ‘Suspend Me Gently’”), но смысл очень бытовой. Если потребитель обрабатывает значения медленно, то хочется, чтобы источник не продолжал выдавать значения бесконечно быстро, иначе мы либо захлебнёмся, либо начнём накапливать значения в памяти.

У Flow (в базовой модели) есть естественный механизм: emit и весь конвейер — suspend-дружелюбны, поэтому если потребитель медленный, то upstream часто будет приостанавливаться и ждать.

Посмотрим на пример: источник быстро эмитит числа, а collect обрабатывает медленно через delay.

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    flow {
        repeat(3) { i ->
            emit(i)
            println("emitted $i") // печать идёт "в такт" обработке
        }
    }.collect { v ->
        delay(100)
        println("collected $v")
        // emitted 0
        // collected 0
        // emitted 1
        // collected 1
        // emitted 2
        // collected 2
    }
}

Если вы ожидали увидеть "emitted 0, emitted 1, emitted 2" сразу, а потом "collected …", то это как раз и есть главный инсайт backpressure в Flow: поток не обязан “выстреливать” всё мгновенно, он может быть естественно синхронизирован с потребителем через приостановки.

7. Собираем конвейер в аккуратный пайплайн

Теперь соберём мини-версию нашего “монитора” так, чтобы код читался как текст: “берём строки → парсим → оставляем важное → логируем → ловим ошибки → печатаем итог”.

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

import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow

fun rawLogLines(): Flow<String> = flow {
    emit("INFO|Service started")
    emit("WARN|Cache miss")
    emit("BROKEN LINE")
    emit("ERROR|Disk is full")
}

Потом опишем конвейер как функцию, которая возвращает Flow<MonitorEvent>. Это хороший стиль: мы отделяем “что делать” от “как выводить”.

import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.filter
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.catch

fun importantEvents(): Flow<MonitorEvent> =
    rawLogLines()
        .map { parseLineStrict(it) }
        .filter { it.level != "INFO" }
        .onEach { println("trace: $it") }
        .catch { e ->
            emit(MonitorEvent("ERROR", "stream failed: ${e.message}"))
        }

И финальный main, где всё собирается:

import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    importantEvents().collect { e ->
        println("ALERT: ${e.level} -> ${e.message}")
    }
}

Обратите внимание, как удобно читать такой код: в importantEvents() нет терминальной операции, значит её можно переиспользовать, тестировать и подключать в разные места. А main остаётся простым: “собери и выведи”.

8. Памятка по операторам

Когда операторов становится больше, полезно держать “карту” в голове. Пусть будет маленькая табличка:

Оператор Что делает Меняет тип потока? Типичный смысл
map { ... }
Преобразует каждое значение Да “Преврати сырое в удобное”
filter { ... }
Пропускает только подходящие Нет “Оставь только нужное”
onEach { ... }
Делает побочный эффект и пропускает дальше Нет “Подсмотри/залоги, не меняя данные”
catch { ... }
Ловит ошибки upstream и даёт альтернативное поведение Обычно нет (но можно emit) “Не падай молча, обработай ошибку”

Если вы видите длинную цепочку, можно буквально читать её вслух: rawLogLines, map, filter, onEach, catch, collect. Обычно это уже хороший знак: код похож на объяснение, а не на загадку.

9. Типичные ошибки

Ошибка №1: делать побочные эффекты в map, а потом удивляться, что “логика какая-то мутная”.
Когда map одновременно и преобразует данные, и печатает, и считает статистику, он перестаёт быть “преобразованием”, а становится “куском всего”. В итоге сложно понять, что именно и зачем меняется. Гораздо спокойнее держать map как чистую трансформацию, а печать и отладку выносить в onEach.

Ошибка №2: поставить catch и считать, что он ловит вообще всё.
catch ловит ошибки, которые произошли до него в цепочке, то есть upstream. Если ошибка возникла внутри collect { ... }, то это уже downstream, и catch может не помочь. В таких ситуациях лучше оборачивать collect в обычный try/catch, как мы делали раньше при изучении исключений.

Ошибка №3: “проглотить” исключение в catch и продолжить как ни в чём не бывало без сигнала.
Пустой catch {} — это как заклеить лампочку “Check Engine” изолентой: машине от этого легче не становится, зато водитель перестаёт понимать, что происходит. Даже если вы выдаёте fallback-значение через emit, полезно хотя бы оставить сообщение об ошибке или явно сформировать “ошибочное событие”, чтобы поведение было предсказуемым.

Ошибка №4: сделать весь конвейер внутри collect и получить “комбайн” на 80 строк.
Такой код сложно тестировать и трудно переиспользовать: логика привязана к месту вывода. Чаще всего лучше вынести конвейер в функцию, которая возвращает Flow<...>, а в main оставить только collect. Это делает код проще и по структуре очень похоже на обычные функции обработки коллекций, только во времени.

Ошибка №5: ожидать, что Flow “настреляет” значения быстрее, чем потребитель их обработает, и не понимать, почему всё идёт медленно.
Это не баг: базовая идея backpressure в том, что Flow умеет естественно подстраиваться под скорость обработки через suspend. Если в collect стоит delay, то upstream часто будет “ждать” — и это как раз помогает не раздувать память и не устраивать лавину событий.

1
Задача
Kotlin SELF, 54 уровень, 3 лекция
Недоступна
Квадраты энергии
Квадраты энергии
1
Задача
Kotlin SELF, 54 уровень, 3 лекция
Недоступна
Фильтр датчика
Фильтр датчика
1
Задача
Kotlin SELF, 54 уровень, 3 лекция
Недоступна
Отладка команд
Отладка команд
1
Задача
Kotlin SELF, 54 уровень, 3 лекция
Недоступна
Парсер логов
Парсер логов
Комментарии
ЧТОБЫ ПОСМОТРЕТЬ ВСЕ КОММЕНТАРИИ ИЛИ ОСТАВИТЬ КОММЕНТАРИЙ,
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ