Sobes.tech
Middle

Kako se prebacuju niti u Flow? Na kojoj niti radi Flow?

sobes.tech АИ

Одговор од АИ

Podrazumevano, Flow radi na niti sa koje je pokrenut. Prebacivanje niti može se izvršiti pomoću operatora flowOn.

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

// Primer funkcije koja vraća Flow
fun simpleFlow(): Flow<Int> = flow {
    logThread("Flow pokrenut") // Logovanje niti gde se pokreće flow builder
    for (i in 1..3) {
        emit(i) // Emitovanje vrednosti
    }
}

fun main() = runBlocking {
    simpleFlow()
        .flowOn(Dispatchers.IO) // Prebacujemo upstream (emittere) na Dispatchers.IO
        .collect { value ->
            logThread("Prikupljeno $value") // Logovanje niti gde se prikupljaju vrednosti
        }
}

// Pomoćna funkcija za logovanje niti
fun logThread(msg: String) {
    println("[$msg] [${Thread.currentThread().name}]")
}

Operator flowOn utiče na nit na kojoj se izvršavaju operatori pre njega u lancu (uključujući flow builder). Operatori posle flowOn se izvršavaju na niti navedenu u argumentu flowOn. Ako u lancu postoji više flowOn, svaki od njih utiče na deo lanca pre njega.

Primer sa više flowOn:

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

fun simpleFlowWithMultipleFlowOn(): Flow<String> = flow {
    logThread("Flow builder pokrenut") // Radiće na prvom specificiranom flowOn
    emit("A")
    emit("B")
}.map {
    logThread("Mapiranje $it") // Radiće na drugom specificiranom flowOn
    it.toLowerCase()
}.flowOn(Dispatchers.Default) // Drugi flowOn utiče na map() i flow builder
.filter {
    logThread("Filtriranje $it") // Radiće na niti sa koje se poziva collect (obično main/runBlocking)
    true
}.flowOn(Dispatchers.IO) // Prvi flowOn utiče na map() i flow builder

fun main() = runBlocking {
    simpleFlowWithMultipleFlowOn().collect { value ->
        logThread("Prikupljanje $value") // Radi na niti runBlocking (main)
    }
}

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

U ovom primeru:

  • flow builder i map će raditi na Dispatchers.Default (jer je to poslednji flowOn u lancu pre njih).
  • filter će raditi na niti sa koje se poziva collect, tj. na niti gde se izvršava runBlocking (obično Main u Androidu).

Važno je razumeti da flowOn kreira međuvremeni buffer i nije "prozirno" prebacivanje.