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 — 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.
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.
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.
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.
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 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.
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.
Első példa — alap Flow számok generálásával és transzformációval a map operátoron keresztül:
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:
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:
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 — 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.
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
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.
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.
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).
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.
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
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.
Olvassa el is