Sobes.tech
Middle

Wie kann man Threads in Flow wechseln? Auf welchem Thread läuft Flow?

sobes.tech KI

Antwort von AI

Standardmäßig läuft Flow auf dem Thread, von dem aus es gestartet wurde. Du kannst Threads mit dem Operator flowOn wechseln.

import kotlinx.coroutines.CoroutineDispatcher // Für CoroutineDispatcher
import kotlinx.coroutines.Dispatchers // Für Dispatchers
import kotlinx.coroutines.flow.Flow // Für Flow
import kotlinx.coroutines.flow.flow // Für flow
import kotlinx.coroutines.flow.flowOn // Für flowOn
import kotlinx.coroutines.runBlocking // Für runBlocking

// Beispiel einer Funktion, die einen Flow zurückgibt
fun simpleFlow(): Flow<Int> = flow {
    logThread("Flow gestartet") // Protokollierung des Threads, auf dem der Flow-Builder läuft
    for (i in 1..3) {
        emit(i) // Werte emittieren
    }
}

fun main() = runBlocking {
    simpleFlow()
        .flowOn(Dispatchers.IO) // Upstream (Emittenten) auf Dispatchers.IO umschalten
        .collect { value ->
            logThread("Gesammelt $value") // Protokollierung des Threads, auf dem die Werte gesammelt werden
        }
}

// Hilfsfunktion zum Protokollieren des Threads
fun logThread(msg: String) {
    println("[$msg] [${Thread.currentThread().name}]")
}

Der Operator flowOn beeinflusst den Thread, auf dem die Operatoren vor ihm in der Kette ausgeführt werden (einschließlich des flow-Builders). Operatoren nach flowOn werden auf dem im Argument flowOn angegebenen Thread ausgeführt. Wenn mehrere flowOn in der Kette vorhanden sind, beeinflusst jeder die Teilkette vor ihm.

Beispiel mit mehreren flowOn:

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

fun simpleFlowWithMultipleFlowOn(): Flow<String> = flow {
    logThread("Flow-Builder gestartet") // Wird auf dem ersten angegebenen flowOn arbeiten
    emit("A")
    emit("B")
}.map {
    logThread("Mapping $it") // Wird auf dem zweiten angegebenen flowOn arbeiten
    it.toLowerCase()
}.flowOn(Dispatchers.Default) // Das zweite flowOn beeinflusst map() und den flow-Builder
.filter {
    logThread("Filtern $it") // Wird auf dem Thread arbeiten, auf dem collect aufgerufen wird (standardmäßig main/runBlocking)
    true
}.flowOn(Dispatchers.IO) // Das erste flowOn beeinflusst map() und den flow-Builder

fun main() = runBlocking {
    simpleFlowWithMultipleFlowOn().collect { value ->
        logThread("Sammeln $value") // Arbeitet auf dem Thread, auf dem runBlocking ausgeführt wird (normalerweise Main)
    }
}

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

In diesem Beispiel:

  • Der flow-Builder und map arbeiten auf Dispatchers.Default (weil dies das letzte flowOn in der Kette vor ihnen ist).
  • filter arbeitet auf dem Thread, auf dem collect aufgerufen wird, also auf dem Thread, auf dem runBlocking ausgeführt wird (normalerweise Main in Android).

Es ist wichtig zu verstehen, dass flowOn einen Zwischenpuffer erstellt und kein "transparenter" Thread-Wechsel ist.