Sobes.tech
Middle

Cum se pot comuta firele în Flow? Pe ce fir rulează Flow?

sobes.tech AI

Răspuns de la AI

În mod implicit, Flow funcționează pe thread-ul de pe care a fost lansat. Poți schimba thread-urile folosind operatorul flowOn.

import kotlinx.coroutines.CoroutineDispatcher // Pentru CoroutineDispatcher
import kotlinx.coroutines.Dispatchers // Pentru Dispatchers
import kotlinx.coroutines.flow.Flow // Pentru Flow
import kotlinx.coroutines.flow.flow // Pentru flow
import kotlinx.coroutines.flow.flowOn // Pentru flowOn
import kotlinx.coroutines.runBlocking // Pentru runBlocking

// Exemplu de funcție care returnează un Flow
fun simpleFlow(): Flow<Int> = flow {
    logThread("Flow a început") // Logare thread unde se inițiază builder-ul de flow
    for (i in 1..3) {
        emit(i) // Emiterea valorilor
    }
}

fun main() = runBlocking {
    simpleFlow()
        .flowOn(Dispatchers.IO) // Schimbăm upstream (emitenți) pe Dispatchers.IO
        .collect { value ->
            logThread("Colectat $value") // Logare thread unde se colectează valorile
        }
}

// Funcție auxiliară pentru logarea thread-ului
fun logThread(msg: String) {
    println("[$msg] [${Thread.currentThread().name}]")
}

Operatorul flowOn influențează thread-ul pe care rulează operatorii înainte de el în lanț (inclusiv builder-ul flow). Operatorii după flowOn rulează pe thread-ul specificat în argumentul flowOn. Dacă există mai multe flowOn în lanț, fiecare influențează partea lanțului înaintea sa.

Exemplu cu mai multe flowOn:

import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.runBlocking

fun simpleFlowWithMultipleFlowOn(): Flow<String> = flow {
    logThread("Builder de Flow început") // Va funcționa pe primul `flowOn` specificat
    emit("A")
    emit("B")
}.map {
    logThread("Mapare $it") // Va funcționa pe al doilea `flowOn` specificat
    it.toLowerCase()
}.flowOn(Dispatchers.Default) // Al doilea `flowOn` influențează map() și builder-ul de flow
.filter {
    logThread("Filtrare $it") // Va funcționa pe thread-ul de unde se apelează `collect` (implicit main/runBlocking)
    true
}.flowOn(Dispatchers.IO) // Primul `flowOn` influențează map() și builder-ul de flow

fun main() = runBlocking {
    simpleFlowWithMultipleFlowOn().collect { value ->
        logThread("Colectare $value") // Funcționează pe thread-ul runBlocking (main)
    }
}

fun logThread(msg: String) {
    println("[$msg] [${Thread.currentThread().name}]")
}

În acest exemplu:

  • Builder-ul de flow și map vor funcționa pe Dispatchers.Default (deoarece acesta este ultimul flowOn din lanț înaintea lor).
  • filter va funcționa pe thread-ul de unde se apelează collect, adică pe thread-ul în care rulează runBlocking (de obicei Main în Android).

Este important să înțelegi că flowOn creează un buffer intermediar și nu reprezintă o "schimbare transparentă de thread".