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 як основне джерело даних (Шар: 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. Якщо виняток виник в операторі після білдера, його перехоплює 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-стан (список завдань, текст пошуку, прапорець завантаження) — кожен підписник отримує актуальне значення. 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 року. Ми проконсультуємо вас і запропонуємо найкраще рішення.
Читайте також