Flow — е тип асинхронен поток от данни от библиотеката Kotlin Coroutines, който имплементира студена семантика. Според Kotlin Documentation, 2025, Flow позволява излъчване на последователност от стойности с операторите map, filter, catch и collect. За разлика от LiveData, Flow е изграден върху корутини и поддържа backpressure.
Основни точки
Flow — е тип от пакета kotlinx.coroutines.flow, който представлява студен асинхронен поток от данни. По същество Flow е корутинна последователност, която излъчва стойности чрез функцията emit() и завършва или успешно, или с изключение. Събирането на потока се извършва чрез терминалния оператор collect(), който е suspend-функция.
Cold stream означава, че кодът вътре в flow-builder-а се изпълнява отново за всеки абонат. Observable.fromIterable в RxJava се държи по подобен начин: нов абонат получава всички стойности от началото. В Flow това е имплементирано чрез suspend-функцията collect, която блокира корутината за цялото време на събиране на данни.
Kotlin предоставя няколко начина за създаване на Flow: flow { } — основна конструкция с emit(), flowOf(vararg values) — за фиксиран набор от стойности, .asFlow() — разширение за колекции и Sequence. Всички builder-и са студени — данните се генерират само при извикване на терминалния оператор.
Разделянето на cold и hot потоци е ключова концепция на реактивното програмиране. Cold stream (Flow, Observable) стартира генериране на данни при абониране. Hot stream (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 ms, което намалява броя на 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-builder-а преди терминалния оператор, то се предава на catch. Ако изключение възникне в оператор след builder-а, catch след този оператор го улавя. retryWhen позволява повторение на абонамента с условие: повтори при мрежова грешка до 3 пъти, но не повтаряй при CancellationException. Flow елиминира грешки, зависими от състояние, тъй като не съхранява състояние — това опростява отстраняването на грешки в сравнение с Observable, където Subject съхранява вътрешно състояние.
Тестването на Flow с kotlinx-coroutines-test използва TestDispatcher за симулация на закъснения. Turbine — популярна библиотека от общността за тестване на Flow: test { } стартира Flow, awaitItem() изчаква следващата стойност, awaitComplete() изчаква завършване. Turbine добавя таймаут по подразбиране, което предотвратява зависването на тестовете. За тестване на StateFlow използвайте .testIn(scope) с проверка на стойностите в хронологичен ред.
Често задавани въпроси
Flow — е асинхронен stream с поддръжка на корутини, оператори и 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 — builder Flow за callback-API. Вътре извикайте registerCallback() с emit(value) в callback-а. 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 г. Ще ви консултираме и ще предложим най-доброто решение.
Прочетете също