Flow — τι είναι, cold και hot streams στις coroutines του Kotlin

Συγγραφέας: IT Sectr Δημοσιεύτηκε: 2026-03-17 Χρόνος ανάγνωσης: 9 λεπ

Flow — είναι ένας τύπος ασύγχρονης ροής δεδομένων από τη βιβλιοθήκη Kotlin Coroutines, που υλοποιεί ψυχρή σημασιολογία. Σύμφωνα με το Kotlin Documentation, 2025, το Flow επιτρέπει την εκπομπή μιας ακολουθίας τιμών με τελεστές map, filter, catch και collect. Σε αντίθεση με το LiveData, το Flow είναι χτισμένο πάνω σε coroutines και υποστηρίζει backpressure.

Κύρια σημεία

  • Flow — ψυχρή ασύγχρονη ροή δεδομένων στο Kotlin Coroutines, δεν εκπέμπει τιμές μέχρι τη συλλογή
  • Cold stream — κάθε συνδρομητής ξεκινά τη δική του ανεξάρτητη εκπομπή από την αρχή
  • Hot stream (SharedFlow, StateFlow) — εκπέμπει τιμές ανεξάρτητα από τους συνδρομητές
  • Τελεστές map, filter, catch, debounce, flatMapLatest μετασχηματίζουν τη ροή χωρίς αποκλεισμό
  • Flow είναι πλήρως συμβατό με το Jetpack Compose μέσω StateFlow και collectAsState()

Τι είναι το Flow στο Kotlin;

Flow — είναι ένας τύπος από το πακέτο kotlinx.coroutines.flow, που αντιπροσωπεύει μια ψυχρή ασύγχρονη ροή δεδομένων. Στην ουσία, το Flow είναι μια ακολουθία coroutine που εκπέμπει τιμές μέσω της συνάρτησης emit() και τελειώνει είτε με επιτυχία είτε με εξαίρεση. Η συλλογή της ροής γίνεται μέσω του τερματικού τελεστή collect(), ο οποίος είναι μια suspend-συνάρτηση.

Ψυχρή σημασιολογία

Cold stream σημαίνει ότι ο κώδικας μέσα στο flow-builder εκτελείται ξανά για κάθε συνδρομητή. Το Observable.fromIterable στο RxJava συμπεριφέρεται παρόμοια: ένας νέος συνδρομητής λαμβάνει όλες τις τιμές από την αρχή. Στο Flow, αυτό υλοποιείται μέσω της suspend-συνάρτησης collect, η οποία μπλοκάρει το coroutine καθ' όλη τη διάρκεια συλλογής δεδομένων.

Flow builders

Το Kotlin παρέχει διάφορους τρόπους δημιουργίας Flow: flow { } — βασική κατασκευή με emit(), flowOf(vararg values) — για σταθερό σύνολο τιμών, .asFlow() — επέκταση για συλλογές και Sequence. Όλοι οι builders είναι ψυχροί — τα δεδομένα παράγονται μόνο κατά την κλήση του τερματικού τελεστή.

Cold και Hot streams

Η διαίρεση σε cold και hot streams είναι μια βασική έννοια του αντιδραστικού προγραμματισμού. Cold stream (Flow, Observable) ξεκινά τη δημιουργία δεδομένων κατά τη συνδρομή. Hot stream (Channel, SharedFlow) εκπέμπει δεδομένα ανεξάρτητα — ο συνδρομητής λαμβάνει μόνο ό,τι συμβαίνει μετά τη συνδρομή, χωρίς την αρχή της ακολουθίας.

SharedFlow — είναι ένα καυτό Flow που μπορεί να έχει πολλούς συνδρομητές και μπορεί να αναπαράγει τις τελευταίες τιμές με ρύθμιση replay. Το SharedFlow είναι κατάλληλο για συμβάντα (εφάπαξ ειδοποιήσεις). StateFlow — η παραλλαγή του με σταθερή τιμή κατάστασης, που αποθηκεύει προσωρινά την τελευταία τιμή για νέους συνδρομητές.

ChannelFlow χρησιμοποιεί Channel στο παρασκήνιο, συνδυάζοντας ιδιότητες του Flow και του Channel. Υποστηρίζει buffering και backpressure μέσω χωρητικότητας (capacity). Το ChannelFlow είναι χρήσιμο κατά τη μετατροπή callback-API σε αντιδραστική ροή, όταν οι τιμές εκπέμπονται από διαφορετικά coroutines.

Μετατροπή μεταξύ cold και hot

Για μετατροπή cold Flow σε hot SharedFlow χρησιμοποιείται ο τελεστής shareIn(scope, started, replay). Η παράμετρος started ελέγχει τη στιγμή εκκίνησης: SharingStarted.WhileSubscribed() — ενεργό όσο υπάρχουν συνδρομητές, Lazily — εκκίνηση στον πρώτο συνδρομητή, Eagerly — άμεση εκκίνηση. Αντίστροφη μετατροπή — hot σε cold: StateFlow.asFlow() επιστρέφει ένα ψυχρό Flow που κατά το collect εκπέμπει την τρέχουσα τιμή του StateFlow. Αυτό είναι βολικό για δοκιμές.

Τελεστές Flow

Flow παρέχει ένα πλούσιο σύνολο τελεστών που λειτουργούν ως suspend-συναρτήσεις μέσα στο coroutine. Οι τελεστές δεν έχουν κατάσταση και επιστρέφουν ένα νέο 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-συναρτήσεις και πρέπει να καλούνται μέσα σε ένα coroutine ή άλλη suspend-συνάρτηση.

Παραδείγματα κώδικα Flow

Πρώτο παράδειγμα — βασικό Flow με δημιουργία αριθμών και μετασχηματισμό μέσω του τελεστή map:

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

scope.launch {
    numberFlow
        .map { "Αριθμός: $it" }
        .collect { value ->
            println(value)
        }
}

Δεύτερο παράδειγμα — μετασχηματισμός ροής με φιλτράρισμα και διαχείριση σφαλμάτων μέσω catch:

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

Τρίτο παράδειγμα — χρήση StateFlow στο ViewModel για αντιδραστικό UI στο 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 και SharedFlow

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 υποστηρίζει coroutines, τελεστές, backpressure και λειτουργεί εκτός του επιπέδου UI. Η μετεγκατάσταση από LiveData σε Flow είναι τυπική πρακτική σε σύγχρονα έργα Android.

Κατά τη χρήση του Flow στο ViewModel, η σωστή επιλογή τύπου είναι σημαντική. StateFlow είναι ιδανικό για κατάσταση UI που πρέπει να επιβιώσει από περιστροφή οθόνης. Το SharedFlow είναι κατάλληλο για συμβάντα όπου η επανεπεξεργασία είναι απαράδεκτη — για παράδειγμα, πλοήγηση. Το Flow με collect() στο lifecycleScope δίνει μέγιστο έλεγχο στο περιβάλλον εκτέλεσης, αλλά απαιτεί χειροκίνητη ακύρωση κατά την έξοδο από την οθόνη.

Η δοκιμή του Flow γίνεται μέσω kotlinx-coroutines-test. Η βιβλιοθήκη παρέχει TestDispatcher — εικονικό χρόνο που επιτρέπει την επιτάχυνση καθυστερήσεων (delay) και τον έλεγχο της σειράς εκτέλεσης coroutines. Το TestScope.runTest { } δημιουργεί απομονωμένο περιβάλλον για δοκιμή Flow. Ο τελεστής toList() χρησιμοποιείται συχνά σε δοκιμές για συλλογή όλων των τιμών flow με timeout, για να ελεγχθεί αν η ροή εξέπεμψε τη σωστή ακολουθία δεδομένων.

Το Flow ενσωματώνεται καλά με το Room (βιβλιοθήκη Android για βάσεις δεδομένων): οι μέθοδοι DAO μπορούν να επιστρέφουν Flow<List<Entity>>. Το Room εκπέμπει αυτόματα μια νέα τιμή σε οποιαδήποτε αλλαγή πίνακα — το UI ενημερώνεται χωρίς χειροκίνητη ενεργοποίηση. Αυτό υλοποιείται μέσω InvalidationTracker, που στο παρασκήνιο χρησιμοποιεί Flow με callbackFlow. Μια τέτοια προσέγγιση εξαλείφει την ανάγκη για LiveData και καθιστά το επίπεδο δεδομένων πλήρως προσανατολισμένο σε coroutines. Το 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 σε coroutines: το έργο μπορεί να ξαναγραφτεί σταδιακά, αφήνοντας μέρος των επιπέδων σε RxJava. Κατά τη μετατροπή πρέπει να ληφθεί υπόψη η διαφορά στη σημασιολογία cold/hot: το Observable μπορεί να είναι και cold και hot, το Flow είναι πάντα cold για κανονικό Flow και hot για SharedFlow.

Διαχείριση σφαλμάτων και δοκιμή Flow

Για τη διαχείριση σφαλμάτων στο 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 και LiveData;

Flow — είναι ένα ασύγχρονο stream με υποστήριξη coroutines, τελεστών και backpressure, που λειτουργεί σε οποιοδήποτε επίπεδο αρχιτεκτονικής. Το LiveData — είναι ένα lifecycle-aware στοιχείο μόνο για το επίπεδο UI. Η Google προτείνει Flow για επιχειρηματική λογική και αποθετήρια, LiveData — για απλές παρατηρήσεις στο ViewModel.

Πότε να χρησιμοποιήσω StateFlow αντί για SharedFlow;

StateFlow — όταν χρειάζεται να αποθηκευτεί κατάσταση UI (λίστα εργασιών, κείμενο αναζήτησης, σημαία φόρτωσης) — κάθε Συνδρομητής λαμβάνει την τρέχουσα τιμή. SharedFlow — για εφάπαξ συμβάντα (πλοήγηση, Snackbar). Το StateFlow δεν πρέπει να χρησιμοποιείται για συμβάντα, καθώς η νέα τιμή μπορεί να υποστεί επεξεργασία ξανά.

Πώς λειτουργεί το backpressure στο Flow;

Στο Flow το backpressure υλοποιείται μέσω μηχανισμού suspend: το emit() σταματά το coroutine αν ο συλλέκτης επεξεργάζεται την προηγούμενη τιμή. Τα κανάλια (Channel) στο ChannelFlow έχουν προσωρινή μνήμη μεγέθους capacity. Σε υπερχείλιση: suspending (αναμονή), drop (απόρριψη) ή conflate (αντικατάσταση με το τελευταίο).

Πώς να μετατρέψω callback σε Flow;

Χρησιμοποιήστε callbackFlow — builder Flow για callback-API. Εσωτερικά καλέστε registerCallback() με emit(value) μέσα στο callback. Το awaitClose εγγυάται κλήση unregisterCallback() κατά την ακύρωση του coroutine. Το callbackFlow υποστηρίζει buffering μέσω Channel(UNLIMITED) στο παρασκήνιο.

Μπορεί το Flow να χρησιμοποιηθεί με RxJava;

Ναι, μέσω μετατροπέων: Flow.asObservable() από το πακέτο kotlinx-coroutines-rx3 μετατρέπει το Flow σε Observable από RxJava 3. Αντίστροφα — CompletableSource.asFlow() για Single/Completable/Maybe. Αυτό είναι χρήσιμο κατά τη μετεγκατάσταση από RxJava σε coroutines σε μεγάλα έργα.

Σύνοψη

  • Flow — ψυχρή ασύγχρονη ροή δεδομένων στο Kotlin Coroutines με suspend-συνάρτηση collect
  • Cold stream ξεκινά την εκπομπή ξανά για κάθε συνδρομητή
  • StateFlow — καυτό δοχείο κατάστασης με προσωρινή αποθήκευση τελευταίας τιμής
  • SharedFlow — καυτή ροή για συμβάντα με ρύθμιση replay και προσωρινής μνήμης
  • Τελεστές map, filter, debounce, catch, flatMapLatest — βάση μετασχηματισμού ροής
  • Google προτείνει το Flow ως κύρια πηγή δεδομένων στη σύγχρονη αρχιτεκτονική Android
  • Το LiveData είναι κατάλληλο μόνο για επίπεδο UI, το Flow — για όλα τα επίπεδα εφαρμογής

Θα αναπτύξουμε μια εφαρμογή για κινητά έτοιμη για χρήση

Η IT Sectr δημιουργεί εφαρμογές iOS και Android για νεοφυείς επιχειρήσεις και επιχειρήσεις από το 2017. Θα σας συμβουλεύσουμε και θα προτείνουμε την καλύτερη λύση.

Συζήτηση έργου

Διαβάστε επίσης