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 подржава корутине, операторе, backpressure и ради изван UI слоја. Миграција са LiveData на Flow је стандардна пракса у савременим Android пројектима.
При коришћењу Flow у ViewModel-у важно је правилно одабрати тип. StateFlow је идеалан за UI стање које треба да преживи ротацију екрана. SharedFlow је погодан за догађаје где је поновна обрада неприхватљива — на пример, навигација. Flow са collect() у lifecycleScope даје максималну контролу над контекстом извршења, али захтева ручно отказивање при изласку из екрана.
Тестирање Flow се врши кроз kotlinx-coroutines-test. Библиотека пружа TestDispatcher — виртуелно време које омогућава убрзавање кашњења (delay) и контролу реда извршења корутина. TestScope.runTest { } ствара изоловано окружење за тестирање Flow. Оператор toList() се често користи у тестовима за сакупљање свих вредности flow са timeout-ом, како би се проверило да је ток емитовао исправну секвенцу података.
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 додаје подразумевани timeout, што спречава зависавање тестова. За тестирање 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 — 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. године. Саветоваћемо вас и предложити најбоље решење.
Прочитајте такође