Введение: вопрос, который задают реже, но ценят больше

«Что будет, если источник генерирует значения быстрее, чем коллектор их обрабатывает?» — вопрос, который почти никто не готовит заранее, и именно он отделяет понимание Flow от заучивания операторов.

Типичный ответ: «Flow умеет справляться с backpressure». А как именно? Что происходит с эмиттером, когда буфер заполнен? Чем conflate отличается от collectLatest? Здесь большинство кандидатов заканчиваются.

Backpressure — это механизм согласования скоростей производителя и потребителя. В RxJava это отдельная большая тема с Flowable и стратегиями, а в Flow она встроена в модель. Разберём, как это работает: базовое поведение, оператор buffer, conflate, collectLatest и горячие потоки. Всё — по официальной документации Kotlin.

Что такое backpressure в Flow

По умолчанию Flow последователен: каждый элемент проходит через все операторы в одной корутине, и эмиттер не может обогнать коллектор. Но эмиссия идёт через внутренний буфер, и когда буфер заполнен, в дело вступает backpressure.

Документация описывает поведение по умолчанию так: коллектор применяет backpressure к источнику — эмиттер приостанавливается, когда буфер полон, и возобновляется, когда коллектор освобождает место. Это стратегия SUSPEND: производитель не теряет данные, он просто ждёт.

Example.kt
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 — выбрасывается значение, которое только что пытались добавить.
Example.kt
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 — отменяет обработку предыдущего значения, когда приходит новое. Разницу видно на примере:

Example.kt
// 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.
TickHandler.kt
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, разберём вашу ситуацию и составим план подготовки.