0Pricing
Kotlin Academy · Lekcja

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

  1. Operatory Flow: map, filter, transform i take
  2. catch i onCompletion: obsługa błędów w Flow
  3. combine i zip: łączenie wielu Flow
  4. flowOn i buffer: kontekst i backpressure
← Powrót do Kotlin Academy