Kotlin Academy · Lezione

flowOn e buffer: contesto e backpressure

Modifichi il contesto di emissione con flowOn e bufferizzi le emissioni per gestire la backpressure.

Lezione 4 di 413 passaggi

flowOn e buffer: contesto e backpressure è una lezione Kotlin Academy gratuita su CoddyKit. Questa è la lezione 4 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento Kotlin Academy, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso Kotlin Academy include 4 lezioni in totale.

Contesto del Flow

Per impostazione predefinita, un flusso viene eseguito nel contesto della coroutine che chiama collect. Il dispatcher del produttore e quello del consumatore coincidono, a meno che non li modifichi.

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 modifica il contesto upstream

flowOn(dispatcher) esegue il flusso upstream (tutto ciò che si trova sopra di esso nella catena) sul dispatcher specificato, mentre la raccolta continua sul dispatcher del chiamante.

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}")
    }
}

Più flowOn in una catena

È possibile utilizzare flowOn più volte. Ogni chiamata influisce sugli operatori direttamente sopra di essa, fino al precedente 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)
}

Il problema del backpressure

Quando il produttore emette valori più velocemente di quanto il consumatore riesca a elaborarli, i valori si accodano. Senza buffering, il produttore deve attendere.

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")
    }
}

Operatore buffer()

buffer() esegue il produttore e il consumatore in modo concorrente, in coroutine separate, memorizzando temporaneamente i valori emessi in un canale. Il produttore non attende il consumatore.

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")
    }
}

Capacità di buffer

buffer(capacity) imposta la dimensione del buffer del canale. Quando il buffer è pieno, il produttore viene sospeso (backpressure). La capacità predefinita è 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() solo per il valore più recente

conflate() scarta i valori intermedi quando il consumatore è lento, mantenendo solo l'emissione più recente. È utile per lo stato dell'interfaccia.

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 per consumatori lenti

collectLatest annulla il blocco di raccolta corrente quando arriva un nuovo valore e lo riavvia con il valore più recente.

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
    }
}

Pattern flowOn + buffer

Combini flowOn e buffer: esegua il produttore su IO (rete/disco), memorizzi temporaneamente i risultati e raccolga i valori su 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 per produttori concorrenti

channelFlow crea un flusso basato su un canale, consentendo a più coroutine di emettere valori in modo concorrente all'interno del builder.

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) }
}

Scelta della strategia corretta

In sintesi: flowOn per cambiare contesto, buffer per aumentare il throughput, conflate per mantenere solo il valore più recente nell'interfaccia, collectLatest per annullare l'elaborazione non più attuale.

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!") }

Controllo rapido

Su cosa influisce flowOn in una catena di flussi?

Riepilogo

flowOn sposta il lavoro upstream su un altro dispatcher. buffer disaccoppia produttore e consumatore per aumentare il throughput. conflate scarta i valori intermedi. collectLatest annulla l'elaborazione lenta quando arrivano nuovi valori.

Gratis per iniziare

Impara Kotlin con un tutor IA — gratis

Scrivi ed esegui vero codice nel tuo browser, ricevi aiuto istantaneo da un tutor IA disponibile 24/7, e riprendi da dove hai lasciato sul web o nell'app.

Corsi
51
Lezioni
203

Domande Frequenti

La lezione «flowOn e buffer: contesto e backpressure» è gratuita?

Sì — il testo completo di «flowOn e buffer: contesto e backpressure» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso Kotlin Academy, passa a CoddyKit PRO. Il corso Kotlin Academy include 4 lezioni in totale.

Cosa imparerò in «flowOn e buffer: contesto e backpressure»?

Modifichi il contesto di emissione con flowOn e bufferizzi le emissioni per gestire la backpressure. Eserciti Kotlin Academy con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.

Ho bisogno di esperienza per iniziare Kotlin Academy?

Non è richiesta alcuna esperienza precedente. Kotlin Academy su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 4 di 4.

Quanto tempo richiede la lezione «flowOn e buffer: contesto e backpressure»?

La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.

Posso scrivere ed eseguire codice in questa lezione Kotlin Academy?

Sì. Ogni lezione Kotlin Academy include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.

Tutte le lezioni di questo corso

  1. Operatori di Flow: map, filter, transform e take
  2. catch e onCompletion: gestione degli errori in Flow
  3. combine e zip: unire più Flow
  4. flowOn e buffer: contesto e backpressure
← Torna a Kotlin Academy