flowOn und buffer: Kontext und Backpressure
Ändern Sie mit flowOn den Emissionskontext und puffern Sie Emissionen zur Backpressure-Behandlung.
flowOn und buffer: Kontext und Backpressure ist eine kostenlose Kotlin Academy-Lektion auf CoddyKit. Dies ist Lektion 4 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des Kotlin Academy-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der Kotlin Academy-Kurs umfasst insgesamt 4 Lektionen.
Flow-Kontext
Standardmäßig wird ein Flow im Kontext der Coroutine ausgeführt, die collect aufruft. Der Dispatcher von Producer und Consumer ist identisch, sofern Sie ihn nicht ändern.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
println("emit on: ${Thread.currentThread().name}")
emit(1)
}.collect {
println("collect on: ${Thread.currentThread().name}")
}
}flowOn ändert den Kontext vorgelagerter Flows
flowOn(dispatcher) führt den vorgelagerten Flow (alles, was in der Kette darüber steht) auf dem angegebenen Dispatcher aus, während die Sammlung auf dem Dispatcher des Aufrufers bleibt.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
println("emit: ${Thread.currentThread().name}")
emit(1)
}.map {
println("map: ${Thread.currentThread().name}")
it * 2
}.flowOn(Dispatchers.Default) // above runs on Default
.collect {
println("collect: ${Thread.currentThread().name}")
}
}Mehrere flowOn-Aufrufe in einer Kette
Sie können flowOn mehrmals verwenden. Jeder Aufruf wirkt auf die Operatoren direkt darüber bis zum vorherigen flowOn.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow { emit(1) }
.map { it + 1 }.flowOn(Dispatchers.IO) // map runs on IO
.map { it * 2 }.flowOn(Dispatchers.Default) // this map runs on Default
.collect { println("Result: $it") } // collect on Main (runBlocking)
}Das Backpressure-Problem
Wenn der Producer schneller Werte emittiert, als der Collector sie verarbeiten kann, sammeln sich die Werte in einer Warteschlange. Ohne Pufferung muss der Producer warten.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
val time = System.currentTimeMillis()
flow {
repeat(3) { i ->
delay(100) // fast producer
emit(i)
}
}.collect {
delay(300) // slow consumer
println("Got $it in ${System.currentTimeMillis() - time}ms")
}
}Der Operator buffer()
buffer() führt Producer und Consumer gleichzeitig in getrennten Coroutines aus und puffert die emittierten Werte in einem Channel. Der Producer muss nicht auf den Consumer warten.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
val time = System.currentTimeMillis()
flow {
repeat(3) { i -> delay(100); emit(i) }
}.buffer() // producer and consumer run concurrently
.collect {
delay(300)
println("Got $it in ${System.currentTimeMillis() - time}ms")
}
}Pufferkapazität von buffer
buffer(capacity) legt die Größe des Channel-Puffers fest. Wenn der Puffer voll ist, wird der Producer suspendiert (Backpressure). Die Standardkapazität beträgt 64.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
flow { repeat(5) { emit(it) } }
.buffer(Channel.RENDEZVOUS) // 0: producer waits
// .buffer(Channel.BUFFERED) // default: 64
// .buffer(Channel.UNLIMITED) // unbounded
.collect { delay(50); println(it) }
}conflate() nur für den neuesten Wert
conflate() verwirft Zwischenwerte, wenn der Collector langsam ist, und behält nur die aktuellste Emission. Das ist für UI-Zustände nützlich.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
repeat(5) { i -> emit(i); delay(50) }
}.conflate()
.collect { i ->
delay(150)
println("Collected: $i") // skips some values
}
}collectLatest für langsame Collector
collectLatest bricht den aktuellen Collection-Block ab, sobald ein neuer Wert eintrifft, und startet ihn mit dem aktuellsten Wert neu.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
flow {
emit(1); delay(50)
emit(2); delay(50)
emit(3)
}.collectLatest { value ->
println("Processing $value")
delay(100) // gets cancelled if new value arrives
println("Done $value") // only prints for last value
}
}Muster mit flowOn + buffer
Kombinieren Sie flowOn und buffer: Führen Sie den Producer auf IO (Netzwerk/Festplatte) aus, puffern Sie die Ergebnisse und sammeln Sie sie auf Main.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun fetchItems(): Flow<String> = flow {
repeat(3) { i ->
delay(100) // simulate IO
emit("item-$i")
}
}.flowOn(Dispatchers.IO).buffer(10)
fun main() = runBlocking {
fetchItems().collect { println("UI: $it") }
}channelFlow für gleichzeitige Producer
channelFlow erstellt einen von einem Channel unterstützten Flow, aus dem mehrere Coroutines innerhalb des Builders gleichzeitig emittieren können.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun concurrentFlow(): Flow<Int> = channelFlow {
launch { send(1) }
launch { send(2) }
launch { send(3) }
}
fun main() = runBlocking {
concurrentFlow().collect { println(it) }
}Die richtige Strategie wählen
Zusammenfassung: flowOn für Kontextwechsel, buffer für höheren Durchsatz, conflate für UI mit ausschließlich dem neuesten Wert und collectLatest zum Abbrechen veralteter Verarbeitung.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
// Guidelines:
// CPU-heavy production -> flowOn(Dispatchers.Default)
// IO-heavy production -> flowOn(Dispatchers.IO)
// Slow consumer -> buffer()
// UI state updates -> conflate() or StateFlow
// Search/autocomplete -> collectLatest or flatMapLatest
fun main() = runBlocking { println("Choose the right strategy!") }Kurzer Test
Worauf wirkt sich flowOn in einer Flow-Kette aus?
Zusammenfassung
flowOn verschiebt vorgelagerte Arbeit auf einen anderen Dispatcher. buffer entkoppelt Producer und Consumer und erhöht so den Durchsatz. conflate verwirft Zwischenwerte. collectLatest bricht langsame Verarbeitung bei neuen Werten ab.
Häufig gestellte Fragen
Ist die Lektion „flowOn und buffer: Kontext und Backpressure“ kostenlos?
Ja — der vollständige Text von „flowOn und buffer: Kontext und Backpressure“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des Kotlin Academy-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der Kotlin Academy-Kurs umfasst insgesamt 4 Lektionen.
Was lerne ich in „flowOn und buffer: Kontext und Backpressure“?
Ändern Sie mit flowOn den Emissionskontext und puffern Sie Emissionen zur Backpressure-Behandlung. Du übst Kotlin Academy mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.
Brauche ich Erfahrung, um Kotlin Academy zu starten?
Keine Vorkenntnisse erforderlich. Kotlin Academy auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 4 von 4.
Wie lange dauert die Lektion „flowOn und buffer: Kontext und Backpressure“?
Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.
Kann ich in dieser Kotlin Academy-Lektion Code schreiben und ausführen?
Ja. Jede Kotlin Academy-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.
Alle Lektionen in diesem Kurs
- Flow-Operatoren: map, filter, transform und take
- catch und onCompletion: Fehlerbehandlung in Flow
- combine und zip: Mehrere Flows zusammenführen
- flowOn und buffer: Kontext und Backpressure