Flow — co to jest, cold i hot strumienie w coroutines Kotlin

Autor: IT Sectr Opublikowano: 2026-03-17 Czas czytania: 9 min

Flow — to typ asynchronicznego strumienia danych z biblioteki Kotlin Coroutines, implementujący zimną semantykę. Według Kotlin Documentation, 2025, Flow pozwala emitować sekwencję wartości z operatorami map, filter, catch i collect. W przeciwieństwie do LiveData, Flow jest zbudowany na korutynach i obsługuje backpressure.

Najważniejsze

  • Flow — zimny asynchroniczny strumień danych w Kotlin Coroutines, nie emituje wartości do momentu kolekcji
  • Cold stream — każdy subskrybent inicjuje własną niezależną emisję od początku
  • Hot stream (SharedFlow, StateFlow) — emituje wartości niezależnie od subskrybentów
  • Operatory map, filter, catch, debounce, flatMapLatest przekształcają strumień bez blokowania
  • Flow jest w pełni kompatybilny z Jetpack Compose przez StateFlow i collectAsState()

Czym jest Flow w Kotlin?

Flow — to typ z pakietu kotlinx.coroutines.flow, reprezentujący zimny asynchroniczny strumień danych. W swej istocie Flow to sekwencja korutynowa, która emituje wartości przez funkcję emit() i kończy się albo sukcesem, albo wyjątkiem. Kolekcja strumienia jest wykonywana przez terminalny operator collect(), który jest funkcją suspend.

Zimna semantyka

Cold stream oznacza, że kod wewnątrz flow-buildera jest wykonywany od nowa dla każdego subskrybenta. Observable.fromIterable w RxJava zachowuje się podobnie: nowy subskrybent otrzymuje wszystkie wartości od początku. W Flow jest to zaimplementowane przez suspend-funkcję collect, która blokuje korutynę na czas zbierania danych.

Flow builders

Kotlin udostępnia kilka sposobów tworzenia Flow: flow { } — podstawowa konstrukcja z emit(), flowOf(vararg values) — dla stałego zestawu wartości, .asFlow() — rozszerzenie dla kolekcji i Sequence. Wszystkie buildery są zimne — dane są generowane tylko przy wywołaniu operatora terminalnego.

Cold i Hot strumienie

Podział na cold i hot strumienie to kluczowa koncepcja programowania reaktywnego. Cold stream (Flow, Observable) uruchamia generowanie danych przy subskrypcji. Hot stream (Channel, SharedFlow) emituje dane niezależnie — subskrybent otrzymuje tylko to, co dzieje się po subskrypcji, bez początku sekwencji.

SharedFlow — to gorący Flow, który może mieć wielu subskrybentów i odtwarzać ostatnie wartości przy ustawieniu replay. SharedFlow nadaje się do zdarzeń (jednorazowe powiadomienia). StateFlow — jego odmiana ze stałą wartością stanu, buforująca ostatnią wartość dla nowych subskrybentów.

ChannelFlow używa Channel pod maską, łącząc właściwości Flow i Channel. Obsługuje buforowanie i backpressure przez pojemność (capacity). ChannelFlow jest przydatny przy konwersji callback-API na strumień reaktywny, gdy wartości są emitowane z różnych korutyn.

Konwersja między cold i hot

Do konwersji cold Flow na hot SharedFlow używa się operatora shareIn(scope, started, replay). Parametr started kontroluje moment uruchomienia: SharingStarted.WhileSubscribed() — aktywny dopóki są subskrybenci, Lazily — uruchomienie przy pierwszym subskrybencie, Eagerly — natychmiastowe uruchomienie. Odwrotna konwersja — hot na cold: StateFlow.asFlow() zwraca zimny Flow, który przy collect emituje bieżącą wartość StateFlow. Jest to wygodne do testowania.

Operatory Flow

Flow udostępnia bogaty zestaw operatorów działających jako suspend-funkcje wewnątrz korutyny. Operatory nie mają stanu i zwracają nowy Flow — oryginalny strumień pozostaje niezmieniony. Pozwala to budować bezpieczne łańcuchy transformacji bez efektów ubocznych.

Operator map przekształca każdą wartość strumienia przez transformację asynchroniczną lub synchroniczną. filter przepuszcza tylko wartości spełniające warunek. catch przechwytuje wyjątki przed operatorem terminalnym i pozwala przywrócić strumień. flatMapLatest anuluje poprzednią emisję przy pojawieniu się nowej wartości — analogicznie do switchMap w Rx.

Operator debounce w Flow opóźnia publikację wartości o określony czas. Jeśli w tym czasie pojawi się nowa wartość — timer jest resetowany. W Androidzie debounce jest używany do wyszukiwania: zapytanie jest wysyłane dopiero po pauzie 300-400 ms, co zmniejsza liczbę wywołań API 3-5 razy.

Operatory terminalne

Oprócz collect(), Flow obsługuje inne operatory terminalne: toList() zbiera wszystkie wartości do listy — przydatne w testach, first() zwraca pierwszy element i anuluje strumień, single() oczekuje dokładnie jednego elementu. fold(initial) akumuluje wartości przez przekazaną funkcję. Wszystkie operatory terminalne to suspend-funkcje i muszą być wywoływane wewnątrz korutyny lub innej suspend-funkcji.

Przykłady kodu Flow

Pierwszy przykład — podstawowy Flow z generowaniem liczb i transformacją przez operator map:

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

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

Drugi przykład — transformacja strumienia z filtracją i obsługą błędów przez catch:

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

Trzeci przykład — użycie StateFlow w ViewModel do reaktywnego UI w 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 — to gorący Flow z pojedynczą bieżącą wartością. Buforuje ostatnią wartość i przekazuje ją nowemu subskrybentowi natychmiast. StateFlow jest kontenerem Observable dla stanu, obsługuje porównanie equals — jeśli nowa wartość pokrywa się z bieżącą, emisja nie następuje. Jetpack Compose używa StateFlow przez collectAsState().

SharedFlow — to bardziej elastyczny gorący Flow bez obowiązkowej wartości początkowej. SharedFlow jest konfigurowany przez replay (liczba wartości dla nowych subskrybentów), extraBufferCapacity (bufor poza replay) i onBufferOverflow (strategia przy przepełnieniu). SharedFlow idealnie nadaje się do zdarzeń jednorazowych: nawigacji, Snackbar, analityki.

Flow w architekturze Android jest zalecany przez Google jako główne źródło danych (Warstwa: Repository → UseCase → ViewModel). LiveData ustępuje Flow pod względem elastyczności: Flow obsługuje korutyny, operatory, backpressure i działa poza warstwą UI. Migracja z LiveData na Flow to standardowa praktyka w nowoczesnych projektach Android.

Przy użyciu Flow w ViewModel ważny jest właściwy wybór typu. StateFlow idealnie nadaje się do stanu UI, który powinien przetrwać obrót ekranu. SharedFlow nadaje się do zdarzeń, gdzie ponowne przetworzenie jest niedopuszczalne — na przykład nawigacja. Flow z collect() w lifecycleScope daje maksymalną kontrolę nad kontekstem wykonania, ale wymaga ręcznej anulacji przy wyjściu z ekranu.

Testowanie Flow wykonuje się przez kotlinx-coroutines-test. Biblioteka udostępnia TestDispatcher — wirtualny czas, który pozwala przyspieszać opóźnienia (delay) i kontrolować kolejność wykonywania korutyn. TestScope.runTest { } tworzy izolowane środowisko do testowania Flow. Operator toList() jest często używany w testach do zebrania wszystkich wartości flow z timeoutem, aby sprawdzić, czy strumień wyemitował prawidłową sekwencję danych.

Flow dobrze integruje się z Room (biblioteka Android dla baz danych): metody DAO mogą zwracać Flow<List<Entity>>. Room automatycznie emituje nową wartość przy każdej zmianie tabeli — UI aktualizuje się bez ręcznego wyzwalacza. Jest to zaimplementowane przez InvalidationTracker, który pod maską używa Flow z callbackFlow. Takie podejście eliminuje potrzebę LiveData i czyni warstwę danych w pełni zorientowaną na korutyny. Jetpack Compose przez collectAsState() subskrybuje StateFlow i przerysowuje tylko te komponenty, których dane się zmieniły — daje to wydajność nieosiągalną z architekturami zorientowanymi na LiveData. DataStore (zamiennik SharedPreferences) również zwraca Flow<Preferences>, zapewniając reaktywne odczytywanie ustawień aplikacji bez ręcznych wyzwalaczy aktualizacji.

Flow obsługuje komunikację międzyprocesową przez kotlinx-coroutines-core na JVM bez dodatkowych bibliotek. Na przykład w aplikacjach serwerowych na Ktor Flow może reprezentować strumień przychodzących wiadomości WebSocket. Każda wiadomość jest emitowana do strumienia, przechodzi filtrację i agregację przez operatory, a wynik jest wysyłany do klienta. Takie podejście zastępuje biblioteki reaktywne takie jak Reactor czy RxJava w projektach Kotlin.

Kompatybilność Flow z istniejącym kodem RxJava zapewnia moduł kotlinx-coroutines-rx3. Funkcja-rozszerzenie Flow.asObservable() konwertuje Flow na Observable z RxJava 3. Odwrotna konwersja — CompletableSource.asFlow(), Observable.asFlow(). Upraszcza to migrację z RxJava na korutyny: można przepisywać projekt etapami, pozostawiając część warstw na RxJava. Przy konwersji należy uwzględnić różnicę w semantyce cold/hot: Observable może być zarówno cold jak i hot, Flow jest zawsze cold dla zwykłego Flow i hot dla SharedFlow.

Obsługa błędów i testowanie Flow

Do obsługi błędów w Flow jest osobliwość: jeśli wyjątek powstał wewnątrz flow-buildera przed operatorem terminalnym, jest przekazywany do catch. Jeśli wyjątek powstał w operatorze po builderze, catch po tym operatorze go przechwytuje. retryWhen pozwala powtórzyć subskrypcję z warunkiem: powtórzyć przy błędzie sieci do 3 razy, ale nie powtarzać przy CancellationException. Flow eliminuje błędy zależne od stanu, ponieważ nie przechowuje stanu — upraszcza to debugowanie w porównaniu z Observable, gdzie Subject przechowuje stan wewnętrzny.

Testowanie Flow z kotlinx-coroutines-test używa TestDispatcher do symulacji opóźnień. Turbine — popularna biblioteka społeczności do testowania Flow: test { } uruchamia Flow, awaitItem() oczekuje następnej wartości, awaitComplete() czeka na zakończenie. Turbine dodaje domyślny timeout, co zapobiega zawieszaniu się testów. Do testowania StateFlow używaj .testIn(scope) z weryfikacją wartości w porządku chronologicznym.

Często zadawane pytania

Jaka jest różnica między Flow a LiveData?

Flow — to asynchroniczny strumień z obsługą korutyn, operatorów i backpressure, działający na każdej warstwie architektury. LiveData — to komponent lifecycle-aware tylko dla warstwy UI. Google zaleca Flow do logiki biznesowej i repozytoriów, LiveData — do prostych obserwacji w ViewModel.

Kiedy używać StateFlow zamiast SharedFlow?

StateFlow — gdy trzeba przechowywać stan UI (lista zadań, tekst wyszukiwania, flagę ładowania) — każdy Subscriber otrzymuje aktualną wartość. SharedFlow — do zdarzeń jednorazowych (nawigacja, Snackbar). StateFlow nie powinien być używany do zdarzeń, ponieważ nowa wartość może być przetworzona ponownie.

Jak działa backpressure w Flow?

W Flow backpressure jest zaimplementowany przez mechanizm suspend: emit() wstrzymuje korutynę, jeśli kolektor przetwarza poprzednią wartość. Kanały (Channel) w ChannelFlow mają bufor o rozmiarze capacity. Przy przepełnieniu: suspending (oczekiwanie), drop (odrzucanie) lub conflate (zastąpienie ostatnią).

Jak przekonwertować callback na Flow?

Użyj callbackFlow — buildera Flow dla callback-API. Wewnątrz wywołaj registerCallback() z emit(value) w callbacku. awaitClose gwarantuje wywołanie unregisterCallback() przy anulacji korutyny. callbackFlow obsługuje buforowanie przez Channel(UNLIMITED) pod maską.

Czy można używać Flow z RxJava?

Tak, przez konwertery: Flow.asObservable() z pakietu kotlinx-coroutines-rx3 przekształca Flow w Observable RxJava 3. Odwrotnie — CompletableSource.asFlow() dla Single/Completable/Maybe. Jest to przydatne przy migracji z RxJava na korutyny w dużych projektach.

Podsumowanie

  • Flow — zimny asynchroniczny strumień danych w Kotlin Coroutines z suspend-funkcją collect
  • Cold stream uruchamia emisję od nowa dla każdego subskrybenta
  • StateFlow — gorący kontener stanu z buforowaniem ostatniej wartości
  • SharedFlow — gorący strumień dla zdarzeń z konfiguracją replay i bufora
  • Operatory map, filter, debounce, catch, flatMapLatest — podstawa transformacji strumienia
  • Google zaleca Flow jako główne źródło danych w nowoczesnej architekturze Android
  • LiveData nadaje się tylko do warstwy UI, Flow — do wszystkich warstw aplikacji

Opracujemy aplikację mobilną pod klucz

IT Sectr tworzy aplikacje na iOS i Androida dla startupów i firm od 2017 roku. Doradzimy Ci i zaproponujemy najlepsze rozwiązanie.

Omów projekt

Przeczytaj również