flowOn и buffer: контекст и противодавление
Изменяйте контекст выдачи с помощью flowOn и буферизуйте выдачу для управления противодавлением.
«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 — локальная установка не требуется.
Все уроки этого курса
- Операторы Flow: map, filter, transform и take
- catch и onCompletion: обработка ошибок в Flow
- combine и zip: объединение нескольких Flow
- flowOn и buffer: контекст и противодавление