flowOn i buffer: kontekst i backpressure
Zmieniaj kontekst emisji za pomocą flowOn i buforuj emisje, aby obsługiwać backpressure.
flowOn i buffer: kontekst i backpressure to bezpłatna lekcja Kotlin Academy na CoddyKit. To lekcja 4 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej Kotlin Academy, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs Kotlin Academy zawiera 4 lekcji w sumie.
Kontekst przepływu
Domyślnie przepływ działa w kontekście korutyny, która wywołuje collect. Dyspozytor producenta i konsumenta jest taki sam, chyba że zostanie zmieniony.
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 zmienia kontekst przepływu źródłowego
flowOn(dispatcher) uruchamia przepływ źródłowy (wszystko powyżej tego operatora w łańcuchu) na określonym dyspozytorze, podczas gdy zbieranie pozostaje na dyspozytorze wywołującego.
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}")
}
}Wiele flowOn w łańcuchu
Można użyć flowOn wiele razy. Każde jego wystąpienie wpływa na operatory znajdujące się bezpośrednio powyżej, aż do poprzedniego 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)
}Problem backpressure
Gdy producent emituje wartości szybciej, niż konsument może je przetwarzać, wartości ustawiają się w kolejce. Bez buforowania producent musi czekać.
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")
}
}Operator buffer()
buffer() uruchamia producenta i konsumenta współbieżnie w oddzielnych korutynach, buforując wyemitowane wartości w kanale. Producent nie czeka na konsumenta.
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")
}
}Pojemność buffer
buffer(capacity) ustawia rozmiar bufora kanału. Gdy bufor jest pełny, producent zostaje zawieszony (backpressure). Domyślna pojemność wynosi 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() tylko z najnowszą wartością
conflate() odrzuca wartości pośrednie, gdy konsument działa wolno, zachowując tylko najnowszą emisję. Jest to przydatne w przypadku stanu interfejsu użytkownika.
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 dla wolnych konsumentów
collectLatest anuluje bieżący blok zbierania, gdy nadejdzie nowa wartość, i uruchamia go ponownie z najnowszą wartością.
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
}
}Wzorzec flowOn + buffer
Należy połączyć flowOn i buffer: uruchamiać producenta na IO (sieć/dysk), buforować wyniki i zbierać je na 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 dla współbieżnych producentów
channelFlow tworzy przepływ oparty na kanale, umożliwiając wielu korutom współbieżne emitowanie wartości wewnątrz konstruktora.
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) }
}Wybór właściwej strategii
Podsumowanie: flowOn służy do przełączania kontekstu, buffer do zwiększania przepustowości, conflate do obsługi tylko najnowszej wartości w interfejsie użytkownika, a collectLatest do anulowania nieaktualnego przetwarzania.
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!") }Szybkie sprawdzenie
Na co wpływa flowOn w łańcuchu przepływu?
Podsumowanie
flowOn przenosi pracę przepływu źródłowego na inny dyspozytor. buffer rozdziela producenta i konsumenta, zwiększając przepustowość. conflate odrzuca wartości pośrednie. collectLatest anuluje wolne przetwarzanie po nadejściu nowych wartości.
Często zadawane pytania
Czy lekcja „flowOn i buffer: kontekst i backpressure” jest bezpłatna?
Tak — pełny tekst „flowOn i buffer: kontekst i backpressure” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu Kotlin Academy, przejdź na CoddyKit PRO. Kurs Kotlin Academy zawiera 4 lekcji w sumie.
Co nauczysz się w „flowOn i buffer: kontekst i backpressure”?
Zmieniaj kontekst emisji za pomocą flowOn i buforuj emisje, aby obsługiwać backpressure. Ćwiczysz Kotlin Academy z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć Kotlin Academy?
Nie wymagamy żadnego doświadczenia. Kotlin Academy w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 4 z 4.
Ile czasu zajmuje lekcja „flowOn i buffer: kontekst i backpressure”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji Kotlin Academy?
Tak. Każda lekcja Kotlin Academy zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Operatory Flow: map, filter, transform i take
- catch i onCompletion: obsługa błędów w Flow
- combine i zip: łączenie wielu Flow
- flowOn i buffer: kontekst i backpressure