0Pricing
Kotlin Academy · Lektion

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

  1. Flow-Operatoren: map, filter, transform und take
  2. catch und onCompletion: Fehlerbehandlung in Flow
  3. combine und zip: Mehrere Flows zusammenführen
  4. flowOn und buffer: Kontext und Backpressure
← Zurück zu Kotlin Academy