Flow — ce este, streamuri cold și hot în corutinele Kotlin

Autor: IT Sectr Publicat: 2026-03-17 Timp de citire: 9 min

Flow — este un tip de flux de date asincron din biblioteca Kotlin Coroutines, care implementează semantica cold. Conform Kotlin Documentation, 2025, Flow permite emiterea unei secvențe de valori cu operatorii map, filter, catch și collect. Spre deosebire de LiveData, Flow este construit pe corutine și suportă backpressure.

Principalele

  • Flow — flux de date asincron cold în Kotlin Coroutines, nu emite valori până la colectare
  • Cold stream — fiecare abonat își inițiază propria emisie independentă de la început
  • Hot stream (SharedFlow, StateFlow) — emite valori independent de abonați
  • Operatorii map, filter, catch, debounce, flatMapLatest transformă fluxul fără blocări
  • Flow este complet compatibil cu Jetpack Compose prin StateFlow și collectAsState()

Ce este Flow în Kotlin?

Flow — este un tip din pachetul kotlinx.coroutines.flow, reprezentând un flux de date asincron cold. În esență, Flow este o secvență corutină care emite valori prin funcția emit() și se termină fie cu succes, fie cu o excepție. Colectarea fluxului se realizează prin operatorul terminal collect(), care este o funcție suspend.

Semantica cold

Cold stream înseamnă că codul din flow-builder se execută din nou pentru fiecare abonat. Observable.fromIterable în RxJava se comportă similar: un nou abonat primește toate valorile de la început. În Flow, acest lucru este implementat prin funcția suspend collect, care blochează corutina pe toată durata colectării datelor.

Flow builders

Kotlin oferă mai multe moduri de a crea Flow: flow { } — construcția de bază cu emit(), flowOf(vararg values) — pentru un set fix de valori, .asFlow() — extensie pentru colecții și Sequence. Toți builderii sunt cold — datele sunt generate doar la apelarea operatorului terminal.

Streamuri Cold și Hot

Împărțirea în streamuri cold și hot este un concept cheie al programării reactive. Cold stream (Flow, Observable) începe generarea datelor la abonare. Hot stream (Channel, SharedFlow) emite date independent — abonatul primește doar ceea ce se întâmplă după abonare, fără începutul secvenței.

SharedFlow — este un Flow hot care poate avea mai mulți abonați și poate reda ultimele valori la setarea replay. SharedFlow este potrivit pentru evenimente (notificări unice). StateFlow — varianta sa cu o valoare de stare fixă, care stochează ultima valoare pentru noii abonați.

ChannelFlow utilizează Channel sub capotă, combinând proprietățile Flow și Channel. Suportă bufferizarea și backpressure prin capacitate (capacity). ChannelFlow este util la conversia callback-API în flux reactive, când valorile sunt emise din diferite corutine.

Conversia între cold și hot

Pentru conversia cold Flow în hot SharedFlow se folosește operatorul shareIn(scope, started, replay). Parametrul started controlează momentul pornirii: SharingStarted.WhileSubscribed() — activ cât există abonați, Lazily — pornire la primul abonat, Eagerly — pornire imediată. Conversia inversă — hot în cold: StateFlow.asFlow() returnează un Flow cold care la collect emite valoarea curentă a StateFlow. Acest lucru este convenabil pentru testare.

Operatorii Flow

Flow oferă un set bogat de operatori care funcționează ca funcții suspend în interiorul corutinei. Operatorii nu au stare și returnează un Flow nou — fluxul original rămâne neschimbat. Acest lucru permite construirea de lanțuri de transformare sigure fără efecte secundare.

Operatorul map transformă fiecare valoare a fluxului printr-o transformare asincronă sau sincronă. filter lasă să treacă doar valorile care îndeplinesc condiția. catch prinde excepțiile înainte de operatorul terminal și permite refacerea fluxului. flatMapLatest anulează emisia anterioară la apariția unei noi valori — similar cu switchMap în Rx.

Operatorul debounce în Flow întârzie publicarea valorii cu un timeout specificat. Dacă în acest timp sosește o nouă valoare — timerul se resetează. În Android, debounce este folosit pentru căutare: cererea este trimisă doar după o pauză de 300-400 ms, ceea ce reduce numărul de apeluri API de 3-5 ori.

Operatorii terminali

Pe lângă collect(), Flow suportă alți operatori terminali: toList() colectează toate valorile într-o listă — util pentru teste, first() returnează primul element și anulează fluxul, single() așteaptă exact un element. fold(initial) acumulează valorile prin funcția transmisă. Toți operatorii terminali sunt funcții suspend și trebuie apelați în interiorul unei corutine sau al altei funcții suspend.

Exemple de cod Flow

Primul exemplu — Flow de bază cu generarea numerelor și transformare prin operatorul map:

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

scope.launch {
    numberFlow
        .map { "Numărul: $it" }
        .collect { value ->
            println(value)
        }
}

Al doilea exemplu — transformarea fluxului cu filtrare și gestionarea erorilor prin catch:

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

Al treilea exemplu — utilizarea StateFlow în ViewModel pentru UI reactive în 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 și SharedFlow

StateFlow — este un Flow hot cu o singură valoare curentă. Stochează ultima valoare și o transmite imediat noului abonat. StateFlow este un container Observable pentru stare, suportă compararea equals — dacă noua valoare coincide cu cea curentă, emisia nu are loc. Jetpack Compose folosește StateFlow prin collectAsState().

SharedFlow — este un Flow hot mai flexibil fără o valoare inițială obligatorie. SharedFlow se configurează prin replay (numărul de valori pentru noii abonați), extraBufferCapacity (buffer dincolo de replay) și onBufferOverflow (strategia la depășire). SharedFlow este ideal pentru evenimente unice: navigare, Snackbar, analitică.

Flow în arhitectura Android este recomandat de Google ca sursă principală de date (Strat: Repository → UseCase → ViewModel). LiveData este inferioară Flow în flexibilitate: Flow suportă corutine, operatori, backpressure și funcționează în afara stratului UI. Migrarea de la LiveData la Flow este o practică standard în proiectele Android moderne.

La utilizarea Flow în ViewModel, alegerea corectă a tipului este importantă. StateFlow este ideal pentru starea UI care trebuie să supraviețuiască rotației ecranului. SharedFlow este potrivit pentru evenimente unde reprocesarea este inacceptabilă — de exemplu, navigarea. Flow cu collect() în lifecycleScope oferă control maxim asupra contextului de execuție, dar necesită anulare manuală la ieșirea din ecran.

Testarea Flow se realizează prin kotlinx-coroutines-test. Biblioteca oferă TestDispatcher — timp virtual care permite accelerarea întârzierilor (delay) și controlul ordinii de execuție a corutinelor. TestScope.runTest { } creează un mediu izolat pentru testarea Flow. Operatorul toList() este adesea folosit în teste pentru a colecta toate valorile flow cu timeout, pentru a verifica dacă fluxul a emis secvența corectă de date.

Flow se integrează bine cu Room (biblioteca Android pentru baze de date): metodele DAO pot returna Flow<List<Entity>>. Room emite automat o nouă valoare la orice modificare a tabelului — UI se actualizează fără declanșator manual. Acest lucru este implementat prin InvalidationTracker, care sub capotă folosește Flow cu callbackFlow. O astfel de abordare elimină necesitatea LiveData și face stratul de date complet orientat pe corutine. Jetpack Compose prin collectAsState() se abonează la StateFlow și redesenează doar componentele ale căror date s-au schimbat — aceasta oferă performanță de neatins cu arhitecturi orientate pe LiveData. DataStore (înlocuitorul SharedPreferences) returnează de asemenea Flow<Preferences>, asigurând citirea reactivă a setărilor aplicației fără declanșatoare manuale de actualizare.

Flow suportă comunicarea interproces prin kotlinx-coroutines-core pe JVM fără biblioteci suplimentare. De exemplu, în aplicații server pe Ktor, Flow poate reprezenta un flux de mesaje WebSocket primite. Fiecare mesaj este emis în flux, trece prin filtrare și agregare prin operatori, iar rezultatul este trimis clientului. O astfel de abordare înlocuiește bibliotecile reactive precum Reactor sau RxJava în proiectele Kotlin.

Compatibilitatea Flow cu codul RxJava existent este asigurată de modulul kotlinx-coroutines-rx3. Funcția extensie Flow.asObservable() convertește Flow în Observable din RxJava 3. Conversia inversă — CompletableSource.asFlow(), Observable.asFlow(). Acest lucru simplifică migrarea de la RxJava la corutine: se poate rescrie proiectul în etape, lăsând o parte din straturi pe RxJava. La conversie trebuie luată în considerare diferența de semantică cold/hot: Observable poate fi atât cold cât și hot, Flow este întotdeauna cold pentru Flow obișnuit și hot pentru SharedFlow.

Gestionarea erorilor și testarea Flow

Pentru gestionarea erorilor în Flow există o particularitate: dacă excepția apare în flow-builder înainte de operatorul terminal, este transmisă în catch. Dacă excepția apare în operatorul după builder, catch după acest operator o prinde. retryWhen permite repetarea abonării cu o condiție: repetă la eroare de rețea de până la 3 ori, dar nu repeta la CancellationException. Flow elimină erorile dependente de stare, deoarece nu stochează stare — acest lucru simplifică depanarea în comparație cu Observable, unde Subject stochează stare internă.

Testarea Flow cu kotlinx-coroutines-test folosește TestDispatcher pentru simularea întârzierilor. Turbine — biblioteca populară din comunitate pentru testarea Flow: test { } pornește Flow, awaitItem() așteaptă următoarea valoare, awaitComplete() așteaptă finalizarea. Turbine adaugă un timeout implicit, ceea ce previne blocarea testelor. Pentru testarea StateFlow folosiți .testIn(scope) cu verificarea valorilor în ordine cronologică.

Întrebări frecvente

Care este diferența dintre Flow și LiveData?

Flow — este un stream asincron cu suport pentru corutine, operatori și backpressure, care funcționează pe orice strat al arhitecturii. LiveData — este o componentă lifecycle-aware doar pentru stratul UI. Google recomandă Flow pentru logica de business și repository-uri, LiveData — pentru observații simple în ViewModel.

Când să folosim StateFlow în loc de SharedFlow?

StateFlow — când trebuie stocată starea UI (lista de sarcini, textul de căutare, flagul de încărcare) — fiecare Abonat primește valoarea actuală. SharedFlow — pentru evenimente unice (navigare, Snackbar). StateFlow nu trebuie folosit pentru evenimente, deoarece noua valoare poate fi procesată din nou.

Cum funcționează backpressure în Flow?

În Flow backpressure este implementat prin mecanismul suspend: emit() oprește corutina dacă colectorul procesează valoarea anterioară. Canalele (Channel) în ChannelFlow au un buffer cu dimensiunea capacity. La depășire: suspending (așteptare), drop (renunțare) sau conflate (înlocuire cu ultima).

Cum să convertim callback în Flow?

Folosiți callbackFlow — builder Flow pentru callback-API. În interior, apelați registerCallback() cu emit(value) în callback. awaitClose garantează apelul unregisterCallback() la anularea corutinei. callbackFlow suportă bufferizarea prin Channel(UNLIMITED) sub capotă.

Se poate folosi Flow cu RxJava?

Da, prin convertoare: Flow.asObservable() din pachetul kotlinx-coroutines-rx3 transformă Flow în Observable RxJava 3. Invers — CompletableSource.asFlow() pentru Single/Completable/Maybe. Acest lucru este util la migrarea de la RxJava la corutine în proiecte mari.

Rezumat

  • Flow — flux de date asincron cold în Kotlin Coroutines cu funcția suspend collect
  • Cold stream pornește emisia din nou pentru fiecare abonat
  • StateFlow — container hot de stare cu stocarea ultimei valori
  • SharedFlow — stream hot pentru evenimente cu configurare replay și buffer
  • Operatorii map, filter, debounce, catch, flatMapLatest — baza transformării fluxului
  • Google recomandă Flow ca sursă principală de date în arhitectura Android modernă
  • LiveData este potrivit doar pentru stratul UI, Flow — pentru toate straturile aplicației

Vom dezvolta o aplicație mobilă la cheie

IT Sectr creează aplicații iOS și Android pentru startup-uri și afaceri din 2017. Vă vom consilia și vă vom propune cea mai bună soluție.

Discutați proiectul

Citiți și