Flow — это тип асинхронного потока данных из библиотеки Kotlin Coroutines, реализующий холодную семантику. По данным Kotlin Documentation, 2025, Flow позволяет эмитить последовательность значений с операторами map, filter, catch и collect. В отличие от LiveData, Flow построен на корутинах и поддерживает backpressure.
Главное
Flow — это тип из пакета kotlinx.coroutines.flow, представляющий холодный асинхронный поток данных. По своей сути Flow — это корутинная последовательность, которая эмитит значения через функцию emit() и завершается либо успешно, либо с исключением. Сбор потока выполняется через терминальный оператор collect(), который является suspend-функцией.
Cold стрим означает, что код внутри flow-билдера выполняется заново для каждого подписчика. Observable.fromIterable в RxJava ведёт себя аналогично: новый подписчик получает все значения с начала. В Flow это реализовано через suspend-функцию collect, которая блокирует корутину на всё время сбора данных.
Kotlin предоставляет несколько способов создания Flow: flow { } — базовая конструкция с emit(), flowOf(vararg values) — для фиксированного набора значений, .asFlow() — расширение для коллекций и Sequence. Все билдеры холодные — данные генерируются только при вызове терминального оператора.
Разделение на cold и hot стримы — ключевая концепция реактивного программирования. Cold стрим (Flow, Observable) запускает генерацию данных при подписке. Hot стрим (Channel, SharedFlow) эмитит данные независимо — подписчик получает только то, что происходит после подписки, без начала последовательности.
SharedFlow — это горячий Flow, который может иметь множество подписчиков и реиграть последние значения при настройке replay. SharedFlow подходит для событий (одноразовые уведомления). StateFlow — его разновидность с фиксированным значением состояния, кэширующим последнее значение для новых подписчиков.
ChannelFlow использует Channel под капотом, сочетая свойства Flow и Channel. Он поддерживает буферизацию и backpressure через ёмкость (capacity). ChannelFlow полезен при конвертации callback-API в реактивный стрим, когда значения эмитятся из разных корутин.
Для конвертации cold Flow в hot SharedFlow используется оператор shareIn(scope, started, replay). Параметр started управляет моментом запуска: SharingStarted.WhileSubscribed() — активен пока есть подписчики, Lazily — запуск при первом подписчике, Eagerly — немедленный запуск. Обратное преобразование — hot в cold: StateFlow.asFlow() возвращает холодный Flow, который при collect эмитит текущее значение StateFlow. Это удобно для тестирования.
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 с генерацией чисел и трансформацией через оператор map:
val numberFlow = flow {
for (i in 1..5) {
delay(500)
emit(i)
}
}
scope.launch {
numberFlow
.map { "Число: $it" }
.collect { value ->
println(value)
}
}
Второй пример — трансформация потока с фильтрацией и обработкой ошибок через catch:
flow {
emit("data1")
emit("data2")
throw RuntimeException("network error")
}
.catch { e ->
emit("fallback_data")
}
.collect { value ->
println(value)
}
Третий пример — использование StateFlow в ViewModel для реактивного UI в Jetpack Compose:
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 — это горячий 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-билдера до терминального оператора, оно пробрасывается в 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 — это асинхронный стрим с поддержкой корутин, операторов и backpressure, работающий на любом слое архитектуры. LiveData — это lifecycle-aware компонент только для UI-слоя. Google рекомендует Flow для бизнес-логики и репозиториев, LiveData — для простых наблюдений в ViewModel.
StateFlow — когда нужно хранить UI-состояние (список задач, текст поиска, флаг загрузки) — каждый Subscriber получает актуальное значение. SharedFlow — для одноразовых событий (навигация, Snackbar). StateFlow не должен использоваться для событий, так как новое значение может быть обработано повторно.
В Flow backpressure реализована через suspend-механизм: emit() приостанавливает корутину, если коллектор обрабатывает предыдущее значение. Каналы (Channel) в ChannelFlow имеют буфер с размером capacity. При переполнении: suspending (ожидание), drop (отбрасывание) или conflate (замена последним).
Используйте callbackFlow — билдер Flow для callback-API. Внутри вызовите registerCallback() с emit(value) внутри колбэка. awaitClose гарантирует вызов unregisterCallback() при отмене корутины. callbackFlow поддерживает буферизацию через Channel(UNLIMITED) под капотом.
Да, через конвертеры: Flow.asObservable() из пакета kotlinx-coroutines-rx3 преобразует Flow в Observable RxJava 3. Обратно — CompletableSource.asFlow() для Single/Completable/Maybe. Это полезно при миграции с RxJava на корутины в крупных проектах.
Итоги
Мы разработаем мобильное приложение под ключ
IT Sectr создаёт приложения для iOS и Android для стартапов и бизнеса с 2017 года. Мы проконсультируем вас и предложим наилучшее решение.
Читайте также