Kotlin Academy · Урок

flowOn и buffer: контекст и противодавление

Изменяйте контекст выдачи с помощью flowOn и буферизуйте выдачу для управления противодавлением.

Урок 4 из 413 шагов

«flowOn и buffer: контекст и противодавление» — бесплатный урок Kotlin Academy на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения Kotlin Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс Kotlin Academy содержит 4 уроков всего.

Контекст потока

По умолчанию поток выполняется в контексте сопрограммы, вызывающей collect. Диспетчеры производителя и получателя совпадают, если явно не изменить это поведение.

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 изменяет контекст исходного потока

flowOn(dispatcher) выполняет исходный поток — всё, что находится выше него в цепочке, — на указанном диспетчере, а сбор значений остаётся на диспетчере вызывающего кода.

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

Несколько flowOn в цепочке

flowOn можно использовать несколько раз. Каждый вызов влияет на операторы, расположенные непосредственно выше него, вплоть до предыдущего 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)
}

Проблема обратного давления

Когда производитель выдаёт значения быстрее, чем сборщик успевает их обрабатывать, значения накапливаются в очереди. Без буферизации это заставляет производителя ждать.

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

Оператор buffer()

buffer() выполняет производителя и сборщика одновременно в отдельных сопрограммах, буферизуя выданные значения в канале. Производителю не приходится ждать сборщика.

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

Ёмкость buffer

buffer(capacity) задаёт размер буфера канала. Когда буфер заполнен, производитель приостанавливается — возникает обратное давление. Значение по умолчанию — 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(): только последнее значение

conflate() отбрасывает промежуточные значения, когда сборщик работает медленно, сохраняя только самое новое значение. Это полезно для состояния интерфейса.

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 для медленных сборщиков

collectLatest отменяет текущий блок сбора, когда поступает новое значение, и запускает его заново с последним значением.

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

Шаблон flowOn + buffer

Объединяйте flowOn и buffer: выполняйте производителя на IO — для сети или диска, — буферизуйте результаты и собирайте их в главном потоке.

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 для параллельных производителей

channelFlow создаёт поток на основе канала, позволяя нескольким сопрограммам одновременно выдавать значения из конструктора.

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

Выбор подходящей стратегии

Итог: flowOn — для переключения контекста, buffer — для повышения пропускной способности, conflate — для интерфейса, которому нужно только последнее значение, collectLatest — для отмены обработки устаревших значений.

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

Быстрая проверка

На что влияет flowOn в цепочке потоков?

Итоги

flowOn переносит работу исходного потока на другой диспетчер. buffer разделяет производителя и сборщика, повышая пропускную способность. conflate отбрасывает промежуточные значения. collectLatest отменяет медленную обработку при поступлении новых значений.

Можно начать бесплатно

Изучай Kotlin с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
51
Уроки
203

Часто задаваемые вопросы

Урок «flowOn и buffer: контекст и противодавление» бесплатный?

Да — полный текст урока «flowOn и buffer: контекст и противодавление» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Kotlin Academy, подпишись на CoddyKit PRO. Курс Kotlin Academy содержит 4 уроков всего.

Чему я научусь в уроке «flowOn и buffer: контекст и противодавление»?

Изменяйте контекст выдачи с помощью flowOn и буферизуйте выдачу для управления противодавлением. Ты практикуешь Kotlin Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать Kotlin Academy?

Предыдущий опыт не требуется. Kotlin Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 4 из 4.

Сколько времени занимает урок «flowOn и buffer: контекст и противодавление»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке Kotlin Academy?

Да. Каждый урок Kotlin Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Операторы Flow: map, filter, transform и take
  2. catch и onCompletion: обработка ошибок в Flow
  3. combine и zip: объединение нескольких Flow
  4. flowOn и buffer: контекст и противодавление
← Назад к Kotlin Academy