Flow — mi ez, cold és hot streamek a Kotlin korutinokban

Szerző: IT Sectr Megjelenés: 2026-03-17 Olvasási idő: 9 perc

Flow — egy aszinkron adatfolyam típus a Kotlin Coroutines könyvtárból, amely hideg szemantikát valósít meg. A Kotlin Documentation, 2025 szerint a Flow lehetővé teszi értékek sorozatának kibocsátását a map, filter, catch és collect operátorokkal. A LiveData-val ellentétben a Flow korutinokra épül és támogatja a backpressure-t.

Főbb pontok

  • Flow — hideg aszinkron adatfolyam a Kotlin Coroutines-ben, nem bocsát ki értékeket a gyűjtésig
  • Cold stream — minden előfizető elindítja saját független kibocsátását az elejétől
  • Hot stream (SharedFlow, StateFlow) — az előfizetőktől függetlenül bocsát ki értékeket
  • Operátorok map, filter, catch, debounce, flatMapLatest blokkolás nélkül alakítják át az adatfolyamot
  • Flow teljes mértékben kompatibilis a Jetpack Compose-zal a StateFlow-n és a collectAsState()-n keresztül

Mi a Flow Kotlinban?

Flow — egy típus a kotlinx.coroutines.flow csomagból, amely egy hideg aszinkron adatfolyamot képvisel. Lényegében a Flow egy korutin sorozat, amely az emit() függvényen keresztül bocsát ki értékeket, és vagy sikeresen vagy kivétellel ér véget. Az adatfolyam gyűjtése a collect() terminális operátoron keresztül történik, amely egy suspend-függvény.

Hideg szemantika

Cold stream azt jelenti, hogy a flow-builderen belüli kód minden előfizető számára újra végrehajtódik. Az Observable.fromIterable az RxJava-ban hasonlóan viselkedik: egy új előfizető az összes értéket az elejétől kapja. A Flow-ban ez a collect suspend-függvényen keresztül van megvalósítva, amely az adatok gyűjtésének teljes időtartamára blokkolja a korutint.

Flow builderek

A Kotlin több módszert kínál a Flow létrehozására: flow { } — alapvető konstrukció emit()-tel, flowOf(vararg values) — rögzített értékkészlethez, .asFlow() — kiterjesztés gyűjteményekhez és Sequence-hez. Minden builder hideg — az adatok csak a terminális operátor meghívásakor jönnek létre.

Cold és Hot streamek

A cold és hot streamekre való felosztás a reaktív programozás egyik kulcsfogalma. Cold stream (Flow, Observable) az adatok generálását az előfizetéskor indítja el. Hot stream (Channel, SharedFlow) az adatokat függetlenül bocsátja ki — az előfizető csak azt kapja, ami az előfizetés után történik, a sorozat eleje nélkül.

SharedFlow — egy forró Flow, amelynek több előfizetője lehet, és a replay beállításakor képes az utolsó értékeket újra lejátszani. A SharedFlow eseményekhez alkalmas (egyszeri értesítések). StateFlow — a változata rögzített állapotértékkel, amely az utolsó értéket gyorsítótárazza az új előfizetők számára.

ChannelFlow a Channel-t használja a motorháztető alatt, egyesítve a Flow és a Channel tulajdonságait. Támogatja a pufferelést és a backpressure-t a kapacitáson (capacity) keresztül. A ChannelFlow hasznos a callback-API reaktív adatfolyammá alakításakor, amikor az értékek különböző korutinokból kerülnek kibocsátásra.

Átalakítás cold és hot között

A cold Flow hot SharedFlow-vá alakításához a shareIn(scope, started, replay) operátort használjuk. A started paraméter szabályozza az indítás pillanatát: SharingStarted.WhileSubscribed() — aktív amíg vannak előfizetők, Lazily — indítás az első előfizetőnél, Eagerly — azonnali indítás. Fordított átalakítás — hot cold-dá: a StateFlow.asFlow() egy hideg Flow-t ad vissza, amely collect-kor a StateFlow aktuális értékét bocsátja ki. Ez kényelmes a teszteléshez.

Flow operátorok

Flow gazdag operátorkészletet biztosít, amelyek suspend-függvényként működnek a korutinon belül. Az operátoroknak nincs állapotuk, és új Flow-t adnak vissza — az eredeti adatfolyam változatlan marad. Ez lehetővé teszi biztonságos transzformációs láncok felépítését mellékhatások nélkül.

A map operátor az adatfolyam minden értékét aszinkron vagy szinkron transzformáción keresztül alakítja át. A filter csak a feltételnek megfelelő értékeket engedi át. A catch elkapja a kivételeket a terminális operátor előtt, és lehetővé teszi az adatfolyam helyreállítását. A flatMapLatest törli az előző kibocsátást egy új érték érkezésekor — hasonlóan a switchMap-hez Rx-ben.

A debounce operátor a Flow-ban egy megadott időtúllépéssel késlelteti az érték közzétételét. Ha ez idő alatt új érték érkezik — az időzítő visszaáll. Androidban a debounce keresésre használatos: a kérés csak 300-400 ms szünet után kerül elküldésre, ami 3-5-ször csökkenti az API-hívások számát.

Terminális operátorok

A collect()-on kívül a Flow más terminális operátorokat is támogat: toList() az összes értéket egy listába gyűjti — hasznos tesztekhez, first() visszaadja az első elemet és törli az adatfolyamot, single() pontosan egy elemet vár. A fold(initial) a megadott függvényen keresztül halmozza fel az értékeket. Minden terminális operátor suspend-függvény, és korutinon vagy más suspend-függvényen belül kell meghívni.

Flow kódpéldák

Első példa — alap Flow számok generálásával és transzformációval a map operátoron keresztül:

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

scope.launch {
    numberFlow
        .map { "Szám: $it" }
        .collect { value ->
            println(value)
        }
}

Második példa — adatfolyam transzformáció szűréssel és hibakezeléssel a catch-en keresztül:

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

Harmadik példa — a StateFlow használata ViewModel-ben reaktív UI-hoz Jetpack Compose-ban:

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 és SharedFlow

StateFlow — egy forró Flow egyetlen aktuális értékkel. Gyorsítótárazza az utolsó értéket, és azonnal továbbítja az új előfizetőnek. A StateFlow egy Observable tároló az állapot számára, támogatja az equals összehasonlítást — ha az új érték megegyezik a jelenlegivel, a kibocsátás nem történik meg. A Jetpack Compose a StateFlow-t a collectAsState()-n keresztül használja.

SharedFlow — egy rugalmasabb forró Flow kötelező kezdeti érték nélkül. A SharedFlow a replay (értékek száma új előfizetők számára), extraBufferCapacity (puffer a replay-en kívül) és onBufferOverflow (stratégia túlcsorduláskor) segítségével konfigurálható. A SharedFlow ideális egyszeri eseményekhez: navigáció, Snackbar, analitika.

Flow az Android architektúrában a Google által ajánlott fő adatforrásként (Réteg: Repository → UseCase → ViewModel). A LiveData rugalmatlanabb a Flow-nál: a Flow támogatja a korutinokat, operátorokat, backpressure-t és az UI rétegen kívül is működik. A LiveData-ról Flow-ra migrálás szabványos gyakorlat a modern Android projektekben.

A Flow ViewModel-ben történő használatakor fontos a megfelelő típus kiválasztása. StateFlow ideális az UI állapothoz, aminek túl kell élnie a képernyő elforgatását. A SharedFlow alkalmas olyan eseményekhez, ahol az újrafeldolgozás elfogadhatatlan — például navigáció. A Flow collect()-tal a lifecycleScope-ban maximális ellenőrzést biztosít a végrehajtási kontextus felett, de kézi törlést igényel a képernyő elhagyásakor.

A Flow tesztelése a kotlinx-coroutines-test segítségével történik. A könyvtár TestDispatcher-t — virtuális időt biztosít, amely lehetővé teszi a késleltetések (delay) gyorsítását és a korutinok végrehajtási sorrendjének szabályozását. A TestScope.runTest { } izolált környezetet hoz létre a Flow teszteléséhez. A toList() operátort gyakran használják tesztekben az összes flow-érték időtúllépéssel történő összegyűjtésére, hogy ellenőrizzék, az adatfolyam a megfelelő adatsorozatot bocsátotta-e ki.

A Flow jól integrálódik a Room-mal (Android könyvtár adatbázisokhoz): a DAO metódusok visszaadhatnak Flow<List<Entity>> értéket. A Room automatikusan új értéket bocsát ki a tábla bármely változásakor — az UI kézi trigger nélkül frissül. Ez az InvalidationTracker segítségével van megvalósítva, amely a motorháztető alatt Flow-t használ a callbackFlow-val. Ez a megközelítés kiküszöböli a LiveData szükségességét, és az adatréteget teljesen korutin-orientáltá teszi. A Jetpack Compose a collectAsState()-n keresztül feliratkozik a StateFlow-ra, és csak azokat a komponenseket rajzolja újra, amelyek adatai megváltoztak — ez olyan teljesítményt nyújt, ami LiveData-orientált architektúrákkal elérhetetlen. A DataStore (a SharedPreferences helyettesítője) szintén Flow<Preferences> értéket ad vissza, biztosítva az alkalmazásbeállítások reaktív olvasását kézi frissítési triggerek nélkül.

A Flow támogatja a folyamatok közötti kommunikációt a kotlinx-coroutines-core segítségével JVM-en, további könyvtárak nélkül. Például a Ktor szerveralkalmazásokban a Flow képviselheti a bejövő WebSocket üzenetek adatfolyamát. Minden üzenet kibocsátásra kerül az adatfolyamba, áthalad a szűrésen és aggregáción az operátorokon keresztül, és az eredmény elküldésre kerül a kliensnek. Ez a megközelítés helyettesíti a reaktív könyvtárakat, mint a Reactor vagy RxJava a Kotlin projektekben.

A Flow kompatibilitását a meglévő RxJava kóddal a kotlinx-coroutines-rx3 modul biztosítja. A Flow.asObservable() kiterjesztési függvény a Flow-t Observable-vé alakítja az RxJava 3-ból. Fordított átalakítás — CompletableSource.asFlow(), Observable.asFlow(). Ez leegyszerűsíti a migrációt RxJava-ról korutinokra: a projekt fokozatosan átírható, a rétegek egy részét RxJava-n hagyva. Az átalakításnál figyelembe kell venni a cold/hot szemantika különbségét: az Observable lehet cold és hot is, a Flow mindig cold a szokásos Flow esetében és hot a SharedFlow esetében.

Hibakezelés és Flow tesztelése

A hibakezelésnek a Flow-ban van egy sajátossága: ha a kivétel a flow-builderen belül, a terminális operátor előtt keletkezik, akkor a catch-be kerül. Ha a kivétel a builder utáni operátorban keletkezik, akkor a catch ezen operátor után fogja el. A retryWhen lehetővé teszi az előfizetés megismétlését feltétellel: ismételd hálózati hiba esetén legfeljebb 3-szor, de ne ismételd CancellationException esetén. A Flow kiküszöböli az állapotfüggő hibákat, mert nem tárol állapotot — ez leegyszerűsíti a hibakeresést az Observable-hez képest, ahol a Subject belső állapotot tárol.

A Flow tesztelése a kotlinx-coroutines-test segítségével a TestDispatcher-t használja a késleltetések szimulálására. A Turbine — népszerű közösségi könyvtár a Flow teszteléséhez: a test { } elindítja a Flow-t, az awaitItem() várja a következő értéket, az awaitComplete() várja a befejezést. A Turbine alapértelmezett időtúllépést ad hozzá, ami megakadályozza a tesztek lefagyását. A StateFlow teszteléséhez használja a .testIn(scope)-t az értékek kronológiai sorrendben történő ellenőrzésével.

Gyakran Ismételt Kérdések

Mi a különbség a Flow és a LiveData között?

Flow — egy aszinkron stream korutinok, operátorok és backpressure támogatásával, amely az architektúra bármely rétegén működik. A LiveData — egy lifecycle-aware komponens csak az UI réteg számára. A Google a Flow-t ajánlja üzleti logikához és repository-khoz, a LiveData-t — egyszerű megfigyelésekhez a ViewModel-ben.

Mikor használjunk StateFlow-t SharedFlow helyett?

StateFlow — amikor UI állapotot kell tárolni (feladatlista, keresőszöveg, betöltési jelző) — minden Előfizető megkapja az aktuális értéket. SharedFlow — egyszeri eseményekhez (navigáció, Snackbar). A StateFlow-t nem szabad eseményekhez használni, mert az új érték újra feldolgozásra kerülhet.

Hogyan működik a backpressure a Flow-ban?

A Flow-ban a backpressure a suspend mechanizmuson keresztül van megvalósítva: az emit() felfüggeszti a korutint, ha a gyűjtő az előző értéket dolgozza fel. A csatornák (Channel) a ChannelFlow-ban capacity méretű pufferrel rendelkeznek. Túlcsorduláskor: suspending (várakozás), drop (eldobás) vagy conflate (helyettesítés az utolsóval).

Hogyan alakítsunk callback-et Flow-vá?

Használja a callbackFlow-t — Flow builder callback-API-hoz. Belül hívja meg a registerCallback()-t az emit(value) függvénnyel a callback-ben. Az awaitClose garantálja az unregisterCallback() meghívását a korutin törlésekor. A callbackFlow támogatja a pufferelést a Channel(UNLIMITED) segítségével a motorháztető alatt.

Használható a Flow RxJava-val?

Igen, konvertereken keresztül: a Flow.asObservable() a kotlinx-coroutines-rx3 csomagból a Flow-t Observable-vé alakítja az RxJava 3-ból. Fordítva — CompletableSource.asFlow() Single/Completable/Maybe számára. Ez hasznos az RxJava-ról korutinokra történő migrációban nagy projektekben.

Összefoglalás

  • Flow — hideg aszinkron adatfolyam a Kotlin Coroutines-ben a collect suspend-függvénnyel
  • Cold stream minden előfizető számára újraindítja a kibocsátást
  • StateFlow — forró állapottároló az utolsó érték gyorsítótárazásával
  • SharedFlow — forró adatfolyam eseményekhez replay és puffer beállításokkal
  • Operátorok map, filter, debounce, catch, flatMapLatest — az adatfolyam-transzformáció alapja
  • Google a Flow-t ajánlja fő adatforrásként a modern Android architektúrában
  • A LiveData csak az UI réteghez alkalmas, a Flow — az alkalmazás minden rétegéhez

Kulcsrakész mobilalkalmazást fejlesztünk

Az IT Sectr 2017 óta készít iOS és Android alkalmazásokat induló vállalkozásoknak és vállalkozásoknak. Tanácsot adunk, és a legjobb megoldást javasoljuk.

Projekt megbeszélése

Olvassa el is