Flow — что это, cold и hot стримы в корутинах Kotlin

Автор: IT Sectr Опубликовано: 2026-03-17 Время чтения: 9 мин

Flow — это тип асинхронного потока данных из библиотеки Kotlin Coroutines, реализующий холодную семантику. По данным Kotlin Documentation, 2025, Flow позволяет эмитить последовательность значений с операторами map, filter, catch и collect. В отличие от LiveData, Flow построен на корутинах и поддерживает backpressure.

Главное

  • Flow — холодный асинхронный поток данных в Kotlin Coroutines, не эмитит значения до сборки
  • Cold стрим — каждый подписчик инициирует свою независимую эмиссию с начала
  • Hot стрим (SharedFlow, StateFlow) — эмитит значения независимо от подписчиков
  • Операторы map, filter, catch, debounce, flatMapLatest трансформируют поток без блокировок
  • Flow полностью совместим с Jetpack Compose через StateFlow и collectAsState()

Что такое Flow в Kotlin?

Flow — это тип из пакета kotlinx.coroutines.flow, представляющий холодный асинхронный поток данных. По своей сути Flow — это корутинная последовательность, которая эмитит значения через функцию emit() и завершается либо успешно, либо с исключением. Сбор потока выполняется через терминальный оператор collect(), который является suspend-функцией.

Холодная семантика

Cold стрим означает, что код внутри flow-билдера выполняется заново для каждого подписчика. Observable.fromIterable в RxJava ведёт себя аналогично: новый подписчик получает все значения с начала. В Flow это реализовано через suspend-функцию collect, которая блокирует корутину на всё время сбора данных.

Flow builders

Kotlin предоставляет несколько способов создания Flow: flow { } — базовая конструкция с emit(), flowOf(vararg values) — для фиксированного набора значений, .asFlow() — расширение для коллекций и Sequence. Все билдеры холодные — данные генерируются только при вызове терминального оператора.

Cold и Hot стримы

Разделение на cold и hot стримы — ключевая концепция реактивного программирования. Cold стрим (Flow, Observable) запускает генерацию данных при подписке. Hot стрим (Channel, SharedFlow) эмитит данные независимо — подписчик получает только то, что происходит после подписки, без начала последовательности.

SharedFlow — это горячий Flow, который может иметь множество подписчиков и реиграть последние значения при настройке replay. SharedFlow подходит для событий (одноразовые уведомления). StateFlow — его разновидность с фиксированным значением состояния, кэширующим последнее значение для новых подписчиков.

ChannelFlow использует Channel под капотом, сочетая свойства Flow и Channel. Он поддерживает буферизацию и backpressure через ёмкость (capacity). ChannelFlow полезен при конвертации callback-API в реактивный стрим, когда значения эмитятся из разных корутин.

Преобразование между cold и hot

Для конвертации cold Flow в hot SharedFlow используется оператор shareIn(scope, started, replay). Параметр started управляет моментом запуска: SharingStarted.WhileSubscribed() — активен пока есть подписчики, Lazily — запуск при первом подписчике, Eagerly — немедленный запуск. Обратное преобразование — hot в cold: StateFlow.asFlow() возвращает холодный Flow, который при collect эмитит текущее значение StateFlow. Это удобно для тестирования.

Операторы Flow

Flow предоставляет богатый набор операторов, работающих как suspend-функции внутри корутины. Операторы не имеют состояния и возвращают новый Flow — оригинальный поток остаётся неизменным. Это позволяет строить безопасные цепочки трансформаций без побочных эффектов.

Оператор map преобразует каждое значение потока через асинхронное или синхронное преобразование. filter пропускает только значения, удовлетворяющие условию. catch перехватывает исключения до терминального оператора и позволяет восстановить поток. flatMapLatest отменяет предыдущую эмиссию при поступлении нового значения — аналогично switchMap в Rx.

Оператор debounce в Flow задерживает публикацию значения на указанный таймаут. Если за это время приходит новое значение — таймер сбрасывается. В Android debounce используется для поиска: запрос уходит только после паузы в 300-400 мс, что сокращает количество вызовов к API в 3-5 раз.

Терминальные операторы

Кроме collect(), Flow поддерживает другие терминальные операторы: toList() собирает все значения в список — полезно для тестов, first() возвращает первый элемент и отменяет поток, single() ожидает ровно один элемент. fold(initial) аккумулирует значения через переданную функцию. Все терминальные операторы — suspend-функции и должны вызываться внутри корутины или другой suspend-функции.

Примеры кода Flow

Первый пример — базовый Flow с генерацией чисел и трансформацией через оператор map:

kotlin
val numberFlow = flow {
    for (i in 1..5) {
        delay(500)
        emit(i)
    }
}

scope.launch {
    numberFlow
        .map { "Число: $it" }
        .collect { value ->
            println(value)
        }
}

Второй пример — трансформация потока с фильтрацией и обработкой ошибок через catch:

kotlin
flow {
    emit("data1")
    emit("data2")
    throw RuntimeException("network error")
}
    .catch { e ->
        emit("fallback_data")
    }
    .collect { value ->
        println(value)
    }

Третий пример — использование StateFlow в ViewModel для реактивного UI в Jetpack Compose:

kotlin
class SearchViewModel : ViewModel() {
    private val _query = MutableStateFlow("")
    val results: StateFlow<List<Result>> = _query
        .debounce(300)
        .flatMapLatest { query ->
            repository.search(query)
        }
        .catch { emit(emptyList()) }
        .stateIn(viewModelScope, SharingStarted.WhileSubscribed(5000), emptyList())

    fun onQueryChanged(query: String) {
        _query.value = query
    }
}

StateFlow и SharedFlow

StateFlow — это горячий Flow с единственным текущим значением. Он кэширует последнее значение и передаёт его новому подписчику немедленно. StateFlow является Observable-контейнером для состояния, поддерживает equals-сравнение — если новое значение совпадает с текущим, эмиссия не происходит. Jetpack Compose использует StateFlow через collectAsState().

SharedFlow — более гибкий горячий Flow без обязательного начального значения. SharedFlow настраивается через replay (количество значений для новых подписчиков), extraBufferCapacity (буфер за пределами replay), и onBufferOverflow (стратегия при переполнении). SharedFlow идеально подходит для одноразовых событий: навигации, Snackbar, аналитики.

Flow в Android-архитектуре рекомендован Google как основной источник данных (Layer: Repository → UseCase → ViewModel). LiveData уступает Flow в гибкости: Flow поддерживает корутины, операторы, backpressure и работает за пределами UI-слоя. Миграция с LiveData на Flow — стандартная практика в современных Android проектах.

При использовании Flow в ViewModel важно правильно выбрать тип. StateFlow идеально подходит для UI-состояния, которое должно переживать поворот экрана. SharedFlow подходит для событий, где повторная обработка недопустима — например, навигация. Flow с collect() в lifecycleScope даёт максимальный контроль над контекстом выполнения, но требует ручной отмены при выходе из экрана.

Тестирование Flow выполняется через kotlinx-coroutines-test. Библиотека предоставляет TestDispatcher — виртуальное время, которое позволяет ускорять задержки (delay) и контролировать порядок выполнения корутин. TestScope.runTest { } создаёт изолированную среду для тестирования Flow. Оператор toList() часто используется в тестах для сбора всех значений flow с таймаутом, чтобы проверить, что поток эмитил правильную последовательность данных.

Flow хорошо интегрируется с Room (библиотека Android для БД): DAO-методы могут возвращать Flow<List<Entity>>. Room автоматически эмитит новое значение при любом изменении таблицы — UI обновляется без ручного триггера. Это реализовано через InvalidationTracker, который под капотом использует Flow с callbackFlow. Такой подход устраняет необходимость в LiveData и делает слой данных полностью корутин-ориентированным. Jetpack Compose через collectAsState() подписывается на StateFlow и перерисовывает только те компоненты, данные которых изменились — это даёт производительность, недостижимую с LiveData-ориентированными архитектурами. DataStore (замена SharedPreferences) также возвращает Flow<Preferences>, обеспечивая реактивное чтение настроек приложения без ручных триггеров обновления.

Flow поддерживает межпроцессное взаимодействие через kotlinx-coroutines-core на JVM без дополнительных библиотек. Например, в серверных приложениях на Ktor Flow может представлять поток входящих WebSocket сообщений. Каждое сообщение эмитится в поток, проходит фильтрацию и агрегацию через операторы, и результат отправляется клиенту. Такой подход заменяет реактивные библиотеки вроде Reactor или RxJava в Kotlin-проектах.

Совместимость Flow с существующим RxJava кодом обеспечивается модулем kotlinx-coroutines-rx3. Функция-расширение Flow.asObservable() конвертирует Flow в Observable из RxJava 3. Обратное преобразование — CompletableSource.asFlow(), Observable.asFlow(). Это упрощает миграцию с RxJava на корутины: можно переписывать проект поэтапно, оставляя часть слоёв на RxJava. При конвертации нужно учитывать разницу в cold/hot семантике: Observable может быть как cold так и hot, Flow всегда cold для обычного Flow и hot для SharedFlow.

Обработка ошибок и тестирование Flow

Для обработки ошибок в Flow есть особенность: если исключение возникло внутри flow-билдера до терминального оператора, оно пробрасывается в catch. Если исключение возникло в операторе после builder, его перехватывает catch после этого оператора. retryWhen позволяет повторить подписку с условием: повторить при network error до 3 раз, но не повторять при CancellationException. Flow исключает state-зависимые ошибки, так как не хранит состояние — это упрощает отладку по сравнению с Observable, где Subject хранит внутреннее состояние.

Тестирование Flow с kotlinx-coroutines-test использует TestDispatcher для симуляции задержек. Turbine — популярная библиотека от сообщества для тестирования Flow: test { } запускает Flow, awaitItem() ожидает следующее значение, awaitComplete() ждёт завершения. Turbine добавляет таймаут по умолчанию, что предотвращает зависание тестов. Для тестирования StateFlow используйте .testIn(scope) с проверкой значений в хронологическом порядке.

Часто задаваемые вопросы

В чём разница между Flow и LiveData?

Flow — это асинхронный стрим с поддержкой корутин, операторов и backpressure, работающий на любом слое архитектуры. LiveData — это lifecycle-aware компонент только для UI-слоя. Google рекомендует Flow для бизнес-логики и репозиториев, LiveData — для простых наблюдений в ViewModel.

Когда использовать StateFlow вместо SharedFlow?

StateFlow — когда нужно хранить UI-состояние (список задач, текст поиска, флаг загрузки) — каждый Subscriber получает актуальное значение. SharedFlow — для одноразовых событий (навигация, Snackbar). StateFlow не должен использоваться для событий, так как новое значение может быть обработано повторно.

Как работает backpressure в Flow?

В Flow backpressure реализована через suspend-механизм: emit() приостанавливает корутину, если коллектор обрабатывает предыдущее значение. Каналы (Channel) в ChannelFlow имеют буфер с размером capacity. При переполнении: suspending (ожидание), drop (отбрасывание) или conflate (замена последним).

Как конвертировать callback в Flow?

Используйте callbackFlow — билдер Flow для callback-API. Внутри вызовите registerCallback() с emit(value) внутри колбэка. awaitClose гарантирует вызов unregisterCallback() при отмене корутины. callbackFlow поддерживает буферизацию через Channel(UNLIMITED) под капотом.

Можно ли использовать Flow с RxJava?

Да, через конвертеры: Flow.asObservable() из пакета kotlinx-coroutines-rx3 преобразует Flow в Observable RxJava 3. Обратно — CompletableSource.asFlow() для Single/Completable/Maybe. Это полезно при миграции с RxJava на корутины в крупных проектах.

Итоги

  • Flow — холодный асинхронный поток данных в Kotlin Coroutines с suspend-функцией collect
  • Cold стрим запускает эмиссию заново для каждого подписчика
  • StateFlow — горячий контейнер состояния с кэшированием последнего значения
  • SharedFlow — горячий стрим для событий с настройкой replay и буфера
  • Операторы map, filter, debounce, catch, flatMapLatest — основа трансформации потока
  • Google рекомендует Flow как основной источник данных в современной Android архитектуре
  • LiveData подходит только для UI-слоя, Flow — для всех слоёв приложения

Мы разработаем мобильное приложение под ключ

IT Sectr создаёт приложения для iOS и Android для стартапов и бизнеса с 2017 года. Мы проконсультируем вас и предложим наилучшее решение.

Обсудить проект

Читайте также