Введение: вопрос, который задают реже, но ценят больше
«Что будет, если источник генерирует значения быстрее, чем коллектор их обрабатывает?» — вопрос, который почти никто не готовит заранее, и именно он отделяет понимание Flow от заучивания операторов.
Типичный ответ: «Flow умеет справляться с backpressure». А как именно? Что происходит с эмиттером, когда буфер заполнен? Чем conflate отличается от collectLatest? Здесь большинство кандидатов заканчиваются.
Backpressure — это механизм согласования скоростей производителя и потребителя. В RxJava это отдельная большая тема с Flowable и стратегиями, а в Flow она встроена в модель. Разберём, как это работает: базовое поведение, оператор buffer, conflate, collectLatest и горячие потоки. Всё — по официальной документации Kotlin.
Что такое backpressure в Flow
По умолчанию Flow последователен: каждый элемент проходит через все операторы в одной корутине, и эмиттер не может обогнать коллектор. Но эмиссия идёт через внутренний буфер, и когда буфер заполнен, в дело вступает backpressure.
Документация описывает поведение по умолчанию так: коллектор применяет backpressure к источнику — эмиттер приостанавливается, когда буфер полон, и возобновляется, когда коллектор освобождает место. Это стратегия SUSPEND: производитель не теряет данные, он просто ждёт.
flow {
repeat(100) { emit(it) } // источник «хочет» отдать всё сразу
}
.collect { value ->
// коллектор обрабатывает медленно
delay(10)
}Ничего не упадёт и ничего не потеряется: эмиттер будет приостанавливаться на emit, пока коллектор не справится с предыдущим значением. Это и есть backpressure — источник не может опередить потребителя, а потребитель не перегружается.
Формула по умолчанию
Flow + коллектор по умолчанию работают по принципу «не обгоняй»: буфер ограничен, при переполнении эмиттер приостанавливается до тех пор, пока коллектор не освободит место. Потерянных значений нет, задержка — есть.
buffer: буфер и политики переполнения
Когда источник должен работать вперёд (например, грузить данные из сети), используется оператор buffer. Он запускает upstream в отдельной корутине и соединяет её с коллектором через канал заданной ёмкости.
По KDoc оператора, buffer(capacity, onBufferOverflow) принимает ёмкость и политику при переполнении:
BufferOverflow.SUSPEND(по умолчанию) — эмиттер ждёт, пока коллектор освободит место. Данные не теряются.BufferOverflow.DROP_OLDEST— при переполнении выбрасывается самое старое значение из буфера.BufferOverflow.DROP_LATEST— выбрасывается значение, которое только что пытались добавить.
import kotlinx.coroutines.channels.BufferOverflow
flow {
repeat(1_000) { emit(it) } // источник работает быстро
}
.buffer(capacity = 4, onBufferOverflow = BufferOverflow.DROP_OLDEST)
.collect { value ->
// обрабатываем — если коллектор не успевает,
// в буфере остаётся только последнее
}Практический смысл: для показа прогресса или позиции плеера старые промежуточные значения не нужны — важнее не отставать от реальности. Для этого и берут DROP_OLDEST или conflate.
Ещё два факта из документации, которые ценятся на собеседовании: buffer(0) убирает буфер полностью — эмиссия и обработка чередуются, а buffer и flowOn при совместном использовании сливаются в один общий буфер (operator fusion), без лишней прослойки.
conflate и collectLatest: только последнее значение
Когда нужен самый свежий результат, а не полная история, есть два инструмента, и их путают чаще всего.
conflate() — это сокращённая запись buffer(1, onBufferOverflow = BufferOverflow.DROP_OLDEST): в буфере всегда одно последнее значение, старые выбрасываются. Важный нюанс из документации: conflate не отменяет обработку, которая уже началась — он влияет только на то, какие значения попадут в буфер.
collectLatest — отменяет обработку предыдущего значения, когда приходит новое. Разницу видно на примере:
// conflate: коллектор всегда дорабатывает текущее значение,
// а пропускает только те, что не успели в буфер
flow { repeat(10) { emit(it) } }
.conflate()
.collect { value -> process(value) }
// collectLatest: обработка старого значения отменяется,
// как только пришло новое
flow { repeat(10) { emit(it) } }
.collectLatest { value ->
process(value) // будет прервана, если придёт следующее значение
}На собеседовании я прошу объяснить разницу, и слышу: «conflate — это чтобы не обрабатывать старые значения». Полуправда. Правильно: conflate выбрасывает из буфера всё, кроме последнего, но не трогает уже начатую обработку; collectLatest отменяет саму обработку. Для долгих операций (парсинг, запись в БД) это принципиальная разница.
Правило выбора: источник генерирует быстрее, чем обрабатывается — и промежуточные значения не важны. UI-сценарий: обновление позиции скролла, громкости, предпросмотра. Тяжёлая обработка — collectLatest, лёгкая — conflate.
Backpressure в горячих потоках: SharedFlow и StateFlow
У горячих потоков своя логика, потому что они не «приостанавливаются по требованию» так же, как холодный Flow. SharedFlow буферизует значения по параметрам конструктора — это описано в документации по Flows:
replay— сколько последних значений получит новый подписчик (кэш повтора).extraBufferCapacity— сколько значений можно хранить сверх replay-кэша, пока коллекторы не успевают.onBufferOverflow— политика, когда буфер полон:SUSPEND,DROP_OLDEST,DROP_LATEST.
class TickHandler(
private val externalScope: CoroutineScope,
private val tickIntervalMs: Long = 5_000
) {
private val _tickFlow = MutableSharedFlow<Unit>(
replay = 0,
extraBufferCapacity = 1,
onBufferOverflow = BufferOverflow.DROP_OLDEST
)
val tickFlow: SharedFlow<Unit> = _tickFlow
init {
externalScope.launch {
while (true) {
_tickFlow.emit(Unit)
delay(tickIntervalMs)
}
}
}
}Здесь эмиттер никогда не приостанавливается: если коллекторы не успевают, лишний «тик» выбрасывается. Без onBufferOverflow emit при полном буфере приостановился бы — а для фонового источника это зависание целой корутины.
StateFlow решает проблему по-другому: он конфлает обновления — всегда хранит одно последнее значение, а медленный коллектор пропускает промежуточные состояния, но всегда получает самое свежее. Параметров буфера у него нет в принципе, и равные по equals значения вообще не считаются обновлением.
Стратегии BufferOverflow
SUSPEND — эмиттер ждёт, DROP_OLDEST — выбрасывается старейшее значение, DROP_LATEST — выбрасывается новое. По умолчанию в Flow и SharedFlow — SUSPEND.
Ошибки в работе с backpressure
- Путают conflate и collectLatest. Первый не трогает начатую обработку, второй отменяет её. Для тяжёлых операций это принципиально.
- Считают, что backpressure в Flow «теряет данные по умолчанию». По умолчанию — SUSPEND: данные не теряются, эмиттер ждёт.
- Не знают, что buffer запускает отдельную корутину. Upstream работает параллельно с коллектором, а не «чуть быстрее в том же потоке».
- Забывают про operator fusion. buffer + flowOn сливаются в один буфер — лишняя прослойка не создаётся.
- Настраивают SharedFlow без onBufferOverflow и удивляются паузам. При полном буфере и политике SUSPEND эмиттер приостанавливается — иногда это не то, что нужно фоновому источнику.
На что я смотрю, когда спрашиваю про backpressure
Я Рустем Бикбулатов, senior Android-разработчик, провожу технические собеседования и менторю в Яндекс Практикуме — больше сотни учеников и разборов за плечами. Вопрос про backpressure я задаю, когда кандидат уверенно прошёл базовые темы: на нём видно, разбирался ли человек в механике Flow или собирал код по примерам.
Сильный ответ: «Flow последователен по умолчанию; buffer позволяет источнику работать вперёд через канал, а при переполнении работает политика — SUSPEND, DROP_OLDEST, DROP_LATEST; conflate — это буфер на одно последнее значение, collectLatest — отмена начатой обработки». Это уровень, когда темы складываются в систему — именно к нему я и веду учеников при подготовке к собеседованиям.
Нужна помощь в переходе на middle?
Бесплатная диагностика: разберу ваши пробелы в Coroutines, Flow и архитектуре и составлю план роста за 30–40 минут. Без обязательств — если менторство вам не подойдёт, честно скажу.
Итоги
Backpressure в Flow — встроенный механизм согласования скоростей, а не отдельная надстройка, как в RxJava.
- По умолчанию: коллектор применяет backpressure — эмиттер приостанавливается при полном буфере, данные не теряются.
- Настройка: buffer с ёмкостью и политикой переполнения; conflate — одно последнее значение; collectLatest — отмена начатой обработки.
- Горячие потоки: SharedFlow настраивается через replay, extraBufferCapacity и onBufferOverflow; StateFlow всегда конфлает до последнего значения.
Главное, что нужно запомнить
Backpressure в Flow — это про то, кто кого ждёт: по умолчанию эмиттер ждёт коллектора, а buffer, conflate и collectLatest позволяют жертвовать промежуточными значениями ради актуальности. Объяснили это — тема закрыта.
Если вопросы про Flow и корутины выбивают вас из колеи на собеседованиях — не нужно зубрить дальше. Напишите в Telegram, разберём вашу ситуацию и составим план подготовки.