Kotlin Academy · Leçon

flowOn et buffer : contexte et contre-pression

Modifiez le contexte d’émission avec flowOn et mettez les émissions en tampon pour gérer la contre-pression.

Leçon 4 sur 413 étapes

flowOn et buffer : contexte et contre-pression est une leçon Kotlin Academy gratuite sur CoddyKit. Ceci est la leçon 4 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage Kotlin Academy, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours Kotlin Academy comprend 4 leçons au total.

Contexte du flux

Par défaut, un flux s’exécute dans le contexte de la coroutine qui appelle collect. Le répartiteur du producteur et celui du consommateur sont identiques, sauf si vous le modifiez.

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 modifie le contexte en amont

flowOn(dispatcher) exécute le flux en amont (tout ce qui le précède dans la chaîne) sur le répartiteur indiqué, tandis que la collecte reste sur le répartiteur de l’appelant.

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

Plusieurs flowOn dans une chaîne

Vous pouvez utiliser flowOn plusieurs fois. Chacune de ses occurrences agit sur les opérateurs qui la précèdent directement, jusqu’à la précédente occurrence de 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)
}

Le problème de la rétropression

Lorsque le producteur émet plus rapidement que le collecteur ne peut traiter les valeurs, celles-ci s’accumulent dans une file d’attente. Sans mise en mémoire tampon, le producteur doit attendre.

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

Opérateur buffer()

buffer() exécute le producteur et le consommateur simultanément dans des coroutines distinctes, en stockant les valeurs émises dans un canal. Le producteur n’attend pas le consommateur.

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é de buffer

buffer(capacity) définit la taille de la mémoire tampon du canal. Lorsque celle-ci est pleine, le producteur est suspendu (rétropression). La capacité par défaut est de 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() pour ne conserver que la plus récente

conflate() supprime les valeurs intermédiaires lorsque le collecteur est lent, en ne conservant que l’émission la plus récente. C’est utile pour l’état de l’interface.

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 pour les collecteurs lents

collectLatest annule le bloc de collecte actuel lorsqu’une nouvelle valeur arrive, puis redémarre avec la valeur la plus récente.

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

Schéma flowOn + buffer

Combinez flowOn et buffer : exécutez le producteur sur IO (réseau/disque), mettez les résultats en mémoire tampon, puis effectuez la collecte sur le répartiteur principal.

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 pour les producteurs concurrents

channelFlow crée un flux reposant sur un canal, ce qui permet à plusieurs coroutines d’émettre simultanément depuis le générateur.

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

Choisir la bonne stratégie

Résumé : flowOn pour changer de contexte, buffer pour le débit, conflate pour ne conserver que la valeur la plus récente dans l’interface, collectLatest pour annuler le traitement obsolète.

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

Vérification rapide

Sur quoi flowOn agit-il dans une chaîne de flux ?

Récapitulatif

flowOn déplace le travail en amont vers un autre répartiteur. buffer découple le producteur et le consommateur pour améliorer le débit. conflate supprime les valeurs intermédiaires. collectLatest annule le traitement lent lorsque de nouvelles valeurs arrivent.

Gratuit pour commencer

Apprends Kotlin avec un tuteur IA — gratuit

Écris et exécute du vrai code dans ton navigateur, obtiens de l'aide instantanée d'un tuteur IA disponible 24h/24, et reprends là où tu t'es arrêté sur le web ou dans l'app.

Cours
51
Leçons
203

Questions Fréquemment Posées

La leçon « flowOn et buffer : contexte et contre-pression » est-elle gratuite ?

Oui — le texte complet de « flowOn et buffer : contexte et contre-pression » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours Kotlin Academy, passe à CoddyKit PRO. Le cours Kotlin Academy comprend 4 leçons au total.

Qu'est-ce que j'apprendrai dans « flowOn et buffer : contexte et contre-pression » ?

Modifiez le contexte d’émission avec flowOn et mettez les émissions en tampon pour gérer la contre-pression. Tu pratiques Kotlin Academy avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.

Dois-je avoir de l'expérience pour commencer Kotlin Academy ?

Aucune expérience préalable n'est requise. Kotlin Academy sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 4 sur 4.

Combien de temps prend la leçon « flowOn et buffer : contexte et contre-pression » ?

La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.

Peux-tu écrire et exécuter du code dans cette leçon Kotlin Academy ?

Oui. Chaque leçon Kotlin Academy inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.

Toutes les leçons de ce cours

  1. Opérateurs de Flow : map, filter, transform et take
  2. catch et onCompletion : gestion des erreurs dans Flow
  3. combine et zip : fusionner plusieurs Flows
  4. flowOn et buffer : contexte et contre-pression
← Retour à Kotlin Academy