Sobes.tech
Middle

Hoe kun je threads in Flow wisselen? Op welke thread werkt Flow?

sobes.tech AI

Antwoord van AI

Standaard werkt Flow op de thread waarop het is gestart. Je kunt van thread wisselen met de operator flowOn.

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

// Voorbeeld van een functie die een Flow teruggeeft
fun simpleFlow(): Flow<Int> = flow {
    logThread("Flow gestart") // Log van de thread waar de flow builder wordt uitgevoerd
    for (i in 1..3) {
        emit(i) // Waarden uitsturen
    }
}

fun main() = runBlocking {
    simpleFlow()
        .flowOn(Dispatchers.IO) // Upstream (emitters) naar Dispatchers.IO
        .collect { value ->
            logThread("Verzameld $value") // Log van de thread waar de waarden worden verzameld
        }
}

// Hulpfunctie voor het loggen van de thread
fun logThread(msg: String) {
    println("[$msg] [${Thread.currentThread().name}]")
}

De operator flowOn beïnvloedt de thread waarop de operators voor hem in de keten worden uitgevoerd (inclusief de flow builder). Operators na flowOn worden uitgevoerd op de thread die in het argument flowOn is gespecificeerd. Als er meerdere flowOn in de keten zijn, beïnvloeden ze elk het deel van de keten dat eraan voorafgaat.

Voorbeeld met meerdere flowOn:

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

fun simpleFlowWithMultipleFlowOn(): Flow<String> = flow {
    logThread("Flow builder gestart") // Werkt op de eerste gespecificeerde flowOn
    emit("A")
    emit("B")
}.map {
    logThread("Mapping $it") // Werkt op de tweede gespecificeerde flowOn
    it.toLowerCase()
}.flowOn(Dispatchers.Default) // De tweede flowOn beïnvloedt map() en flow builder
.filter {
    logThread("Filtering $it") // Werkt op de thread waar `collect` wordt aangeroepen (standaard main/runBlocking)
    true
}.flowOn(Dispatchers.IO) // De eerste flowOn beïnvloedt map() en flow builder

fun main() = runBlocking {
    simpleFlowWithMultipleFlowOn().collect { value ->
        logThread("Verzameld $value") // Werkt op de thread van runBlocking (main)
    }
}

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

In dit voorbeeld:

  • De flow builder en map werken op Dispatchers.Default (omdat dit de laatste flowOn in de keten voor hen is).
  • filter werkt op de thread waarop collect wordt aangeroepen, dus op de thread waar runBlocking wordt uitgevoerd (meestal Main in Android).

Het is belangrijk te begrijpen dat flowOn een tussenbuffer creëert en geen "transparante" thread-switch is.