Senior
Kotlin Flow: подготовка к Middle+/Senior собеседованию
Подробная лекция о Kotlin Flow: cold и hot потоках, StateFlow, SharedFlow, операторах, Android UI и вопросах уровня Senior.
На этой странице
Конспект по Kotlin Flow: от базовой модели до типичных Senior-вопросов. Все исходные примеры сохранены как read-only: они не запускались при подготовке, потому что исходник не задаёт автономную конфигурацию зависимостей.
Теория
Что такое Flow
Flow<T> — асинхронный поток значений.
suspend fun getUser(): User
возвращает одно значение в будущем, а:
fun observeUsers(): Flow<List<User>>
может выдавать новые значения много раз:
time ──────────────────────────────>
users1 users2 users3
Flow ───●──────────●──────────●─────>
Типичный Android pipeline:
Room / DataStore / Repository
↓
Flow
↓
ViewModel
↓
StateFlow<UiState>
↓
collectAsStateWithLifecycle()
↓
Compose
Producer → operators → collector
Flow удобно мыслить как pipeline:
Producer
│
▼
flow { emit(...) }
│
▼
map
│
▼
filter
│
▼
debounce
│
▼
Collector
Пример:
repository.observeUsers()
.map { users -> users.filter(User::isActive) }
.distinctUntilChanged()
.collect { users ->
render(users)
}
Здесь:
observeUsers()— источник;map,distinctUntilChanged— intermediate operators;collect— terminal operator.
Cold Flow
Обычный Flow чаще всего cold.
val numbers = flow {
println("Started")
emit(1)
emit(2)
emit(3)
}
Создание Flow ничего не запускает.
val numbers = flow { ... }
Producer запускается только при collection:
numbers.collect {
println(it)
}
Если подписаться дважды:
numbers.collect { println("A: $it") }
numbers.collect { println("B: $it") }
upstream будет запущен два раза.
Collector A ──► producer execution #1
Collector B ──► producer execution #2
Формулировка для собеседования:
Cold Flow не существует как независимо работающий producer: каждый collector запускает upstream заново.
Создание Flow
Основной builder:
fun numbers(): Flow<Int> = flow {
emit(1)
emit(2)
emit(3)
}
Внутри можно вызывать suspend-функции:
fun users(): Flow<User> = flow {
val users = api.loadUsers()
users.forEach {
emit(it)
}
}
Другие способы:
flowOf(1, 2, 3)
listOf(1, 2, 3).asFlow()
Terminal operators и collect
Пока terminal operator не вызван, cold pipeline обычно не выполняется.
val result = repository.users()
.map { ... }
.filter { ... }
Запуск:
result.collect {
println(it)
}
Другие terminal operators:
first()
firstOrNull()
single()
singleOrNull()
toList()
reduce()
fold()
launchIn(scope)
collect — suspend-функция
viewModelScope.launch {
repository.users().collect {
println(it)
}
println("DONE")
}
Если Flow бесконечный, до DONE нормальное выполнение не дойдёт.
Два бесконечных collect подряд
Ошибка:
viewModelScope.launch {
flowA.collect { ... }
flowB.collect { ... }
}
Если flowA бесконечный, второй collect никогда не стартует.
Можно использовать отдельные coroutines:
viewModelScope.launch {
flowA.collect { ... }
}
viewModelScope.launch {
flowB.collect { ... }
}
или декларативно объединить источники через combine.
Последовательность выполнения
Flow sequential по умолчанию:
flowOf(1, 2, 3)
.map {
println("map $it")
it * 2
}
.collect {
println("collect $it")
}
Модель:
emit(1)
↓
map(1)
↓
collect(2)
↓
emit(2)
↓
map(2)
↓
collect(4)
То есть producer и downstream по умолчанию связаны suspension.
Backpressure, buffer, conflate
Backpressure
flow {
repeat(1000) {
emit(it)
}
}.collect {
delay(1000)
}
Producer не обязан складывать 1000 элементов в память. emit() может приостановиться, пока downstream не будет готов продолжить.
Упрощённо:
Producer Collector
emit(1) ──────────────► process 1
│ │
│ suspend │
│ │
emit(2) ◄──────────────────┘
buffer()
flow
.buffer()
.collect {
slowOperation(it)
}
buffer() позволяет upstream и downstream работать более независимо.
Producer Buffer Consumer
● ─────────► [●]
● ─────────► [● ●] ──────► slow
● ─────────► [● ●]
Он также создаёт coroutine boundary между частями pipeline.
conflate()
Если промежуточные значения не важны:
flow
.conflate()
.collect {
render(it)
}
Например:
Producer: 1 2 3 4 5 6 7 8
Collector: 1 ───── 4 ───── 8
Collector получает актуальные значения, но может пропускать промежуточные.
Latest-операторы
collectLatest
flow.collectLatest { value ->
process(value)
}
Если приходит новое значение, обработка предыдущего отменяется.
"a"
└── process("a") ─────X
"ab"
└── process("ab") ────X
"abc"
└── process("abc") ───── DONE
mapLatest
query
.mapLatest { query ->
repository.search(query)
}
Новое значение отменяет transform предыдущего.
Важное отличие
conflate() позволяет пропускать промежуточные значения.
collectLatest() отменяет работу collector, запущенную для предыдущего значения.
debounce и distinctUntilChanged
debounce
query.debounce(300)
Значение проходит дальше, если в течение указанного интервала не появилось более новое.
A──50ms──Al──80ms──Ale──100ms──Alex
│
│ 300 ms
▼
emit
Типичный поиск:
query
.debounce(300)
.filter { it.length >= 2 }
.distinctUntilChanged()
.mapLatest(repository::search)
distinctUntilChanged
1 1 1 2 2 3 3 2
становится:
1 2 3 2
Удаляются только последовательные одинаковые значения.
Ошибки и cancellation
catch
flow {
emit(api.loadUsers())
}
.catch {
emit(emptyList())
}
catch ловит ошибки upstream относительно своего положения.
flow
.catch { ... }
.collect {
throw RuntimeException()
}
Этот exception из тела collect данный catch не поймает.
upstream
│
▼
catch
│
▼
collector
│
X exception
Cancellation
Flow построен на coroutines:
val job = scope.launch {
flow.collect { ... }
}
job.cancel()
Collection будет отменён.
CancellationException
Нельзя бездумно проглатывать cancellation:
try {
...
} catch (e: CancellationException) {
throw e
} catch (e: Throwable) {
...
}
Отмена coroutine — часть нормального control flow.
CoroutineContext и flowOn
Flow соблюдает context preservation.
Неправильная попытка переключить dispatcher:
flow {
withContext(Dispatchers.IO) {
emit(loadData())
}
}
Для изменения контекста upstream используется:
flowOn()
Пример:
repository.users()
.map(::heavyTransformation)
.flowOn(Dispatchers.Default)
.collect {
render(it)
}
Главное правило:
flowOnвлияет на upstream относительно своего положения, а не на downstream.
flow
.map { A(it) }
.flowOn(Dispatchers.IO)
.map { B(it) }
.collect { C(it) }
Если collector работает на Main:
Dispatchers.IO:
producer
↓
A()
---------- flowOn boundary ----------
Main:
B()
↓
C()
Это частый Senior-вопрос.
combine и zip
zip
a.zip(b) { x, y ->
x to y
}
Сопоставляет элементы попарно:
A: A1 ─── A2 ─── A3
B: B1 ─────── B2 ─── B3
A1B1 A2B2 A3B3
combine
combine(
userFlow,
settingsFlow,
) { user, settings ->
UiState(user, settings)
}
После появления первого значения у каждого source, изменение любого source пересчитывает результат с использованием latest values.
User:
U1────────────U2────────>
Settings:
S1──S2─────────────>
combine:
U1S1
U1S2
U2S2
На собеседовании
zipсвязывает элементы попарно.combineмоделирует реактивное состояние из последних значений нескольких источников.
Для Android UI значительно чаще нужен combine.
flatMapConcat / Merge / Latest
Допустим:
val ids: Flow<UserId>
fun observeUser(id: UserId): Flow<User>
Если сделать:
ids.map {
observeUser(it)
}
получим:
Flow<Flow<User>>
flatMapConcat
Inner flows выполняются последовательно.
A → [1 2 3]
B → [4 5 6]
C → [7 8]
flatMapMerge
Inner flows могут выполняться конкурентно:
A → [────────]
B → [────────]
C → [────────]
Результаты merge'ятся по мере готовности.
flatMapLatest
selectedUserId
.flatMapLatest(repository::observeUser)
Новое outer value отменяет предыдущий inner Flow.
User 1: [────────────X]
User 2: [──────────X]
User 3: [────────────>
Частые сценарии:
- selected user;
- selected account;
- selected chat;
- selected venue;
- search query.
Hot Flow
SharedFlow и StateFlow — hot streams.
┌── Collector A
Producer ──────┼── Collector B
└── Collector C
Один hot stream может существовать независимо от конкретного collector.
Основные реализации:
StateFlow<T>
SharedFlow<T>
StateFlow
private val _state =
MutableStateFlow(UiState())
val state: StateFlow<UiState> =
_state.asStateFlow()
Изменение:
_state.value = newState
Для read-modify-write:
_state.update {
it.copy(isLoading = true)
}
StateFlow всегда имеет текущее значение
state.value
доступно синхронно.
Новый collector сразу получает текущее состояние.
current state = S3
new collector
│
▼
S3
Equality-based conflation
Если новое значение equals() старому, повторная emission подавляется.
Поэтому хороший state:
data class UiState(
val loading: Boolean = false,
val users: List<User> = emptyList(),
)
State должен быть immutable
Плохо:
StateFlow<MutableList<User>>
Изменение внутреннего mutable объекта может не создать новую emission.
Лучше:
_state.update {
it.copy(users = it.users + user)
}
Почему update, а не _state.value = ...
Потенциально проблемно:
_state.value = _state.value.copy(
counter = _state.value.counter + 1
)
При конкурентном read-modify-write update может потеряться.
Предпочтительно:
_state.update {
it.copy(counter = it.counter + 1)
}
SharedFlow
private val _events =
MutableSharedFlow<Event>()
val events =
_events.asSharedFlow()
Emission:
_events.emit(Event.LoggedOut)
replay
MutableSharedFlow<Event>(
replay = 1
)
A → B → C
new subscriber:
replay = 0 → ничего
replay = 1 → C
replay = 2 → B, C
StateFlow vs SharedFlow
Удобная mental model:
StateFlow:
"Что сейчас?"
SharedFlow:
"Какие значения broadcast'ятся?"
StateFlow специализирован под state и имеет:
- обязательный initial value;
.value;- equality-based conflation.
SharedFlow более общий и позволяет управлять:
replay
extraBufferCapacity
onBufferOverflow
emit vs tryEmit
_events.emit(event)
— suspend-функция.
val accepted = _events.tryEmit(event)
— попытка emission без suspension.
Поведение зависит от replay, buffer и наличия subscribers.
Buffer
MutableSharedFlow<Event>(
replay = 0,
extraBufferCapacity = 64,
onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
Overflow policies:
BufferOverflow.SUSPEND
BufferOverflow.DROP_OLDEST
BufferOverflow.DROP_LATEST
shareIn и stateIn
shareIn
Cold:
Collector A ──► producer A
Collector B ──► producer B
После shareIn:
val shared = repository.users()
.shareIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(),
replay = 1,
)
┌── Collector A
producer ─ shared ┤
└── Collector B
stateIn
val state: StateFlow<UiState> =
repository.observeUsers()
.map { users ->
UiState(users = users)
}
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5_000),
initialValue = UiState(),
)
SharingStarted
Eagerly
Upstream запускается сразу.
Lazily
Upstream запускается после появления первого subscriber.
WhileSubscribed
Upstream активен в зависимости от наличия subscribers.
Частый Android-вариант:
SharingStarted.WhileSubscribed(5_000)
Timeout помогает переживать короткие промежутки без подписчика, например lifecycle/configuration transitions, не перезапуская upstream мгновенно.
Flow в Android и Compose
Типичная ViewModel:
class UsersViewModel(
repository: UserRepository,
) : ViewModel() {
val uiState: StateFlow<UsersUiState> =
combine(
repository.observeUsers(),
repository.observePermissions(),
) { users, permissions ->
UsersUiState(
users = users,
canEdit = permissions.canEdit,
)
}
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5_000),
initialValue = UsersUiState(),
)
}
В Compose:
val uiState by viewModel.uiState
.collectAsStateWithLifecycle()
Pipeline:
Repository
↓
Flow
↓
ViewModel
↓
StateFlow<UiState>
↓
collectAsStateWithLifecycle()
↓
Compose State
↓
Recomposition
Для View-system:
lifecycleScope.launch {
repeatOnLifecycle(Lifecycle.State.STARTED) {
viewModel.state.collect {
render(it)
}
}
}
callbackFlow
Нужен для адаптации callback API.
fun locations(): Flow<Location> =
callbackFlow {
val listener = object : Listener {
override fun onLocation(location: Location) {
trySend(location)
}
}
locationManager.register(listener)
awaitClose {
locationManager.unregister(listener)
}
}
Модель:
Callback API
↓
callback
↓
trySend()
↓
Channel
↓
Flow
↓
Collector
awaitClose
При cancellation collector:
collector cancelled
↓
callbackFlow closes
↓
awaitClose
↓
unregister(listener)
Без cleanup возможны:
- resource leaks;
- listener leaks;
- duplicate callbacks.
Flow vs Channel / Sequence / suspend
Flow vs Channel
Упрощённо:
Flow
→ stream abstraction / declarative pipeline
Channel
→ coroutine communication primitive
Некоторые Flow operators/builders используют channels внутри:
callbackFlow;channelFlow;buffer;flowOn.
Но Flow — не просто alias для Channel.
Flow vs Sequence
Оба могут быть lazy.
Но:
Sequence → synchronous
Flow → asynchronous + suspend
В Flow можно:
flow {
delay(1000)
emit(1)
}
Flow vs suspend function
Если операция возвращает один результат:
suspend fun login(): User
Если это поток изменений:
fun observeUser(): Flow<User>
Не стоит оборачивать одноразовую операцию в Flow без причины.
retry и retryWhen
repository.load()
.retry(3) {
it is IOException
}
Важно:
retryперезапускает upstream.
Если upstream содержит non-idempotent side effect:
flow {
chargeCreditCard()
emit(...)
}
слепой retry опасен.
retryWhen
.retryWhen { cause, attempt ->
if (
cause is IOException &&
attempt < 3
) {
delay(1000L * (attempt + 1))
true
} else {
false
}
}
Позволяет реализовать conditional retry и backoff.
Типичные архитектурные ошибки
Два бесконечных collect подряд
flowA.collect { ... }
flowB.collect { ... }
Второй может никогда не запуститься.
Flow для одноразовой операции
fun login(): Flow<User>
когда достаточно:
suspend fun login(): User
Mutable state внутри StateFlow
StateFlow<MutableList<User>>
Усложняет контроль изменений и equality semantics.
Использовать SharedFlow автоматически для каждого UI event
Нужно сначала спросить:
Это действительно transient event или часть состояния?
Если UI обязан восстановить информацию после отсутствия collector, state часто надёжнее transient event.
Несколько collectors дорогого cold upstream
expensiveFlow.collect(...)
expensiveFlow.collect(...)
могут запустить upstream дважды.
Рассмотреть:
shareIn(...)
stateIn(...)
Неправильное понимание flowOn
flowOn не означает «всё после него работает на этом dispatcher».
Он влияет на upstream.
Проглатывание cancellation
Broad catch не должен ломать нормальную отмену coroutine.
Senior-механика Flow
emit — suspend
emit() — не обычный callback.
flow {
emit(value)
}
Он может приостанавливать producer из-за downstream/backpressure.
Отсюда вытекает важное свойство обычного pipeline:
upstream и downstream
координируются через suspension
Context preservation
Обычный flow {} не должен произвольно emit'ить из другого coroutine context.
Это позволяет collector рассуждать о pipeline без скрытых context jumps.
Для смены upstream dispatcher:
flowOn(Dispatchers.IO)
Для concurrent producers существуют другие primitives, например channelFlow.
Exception transparency
Flow старается сохранять прозрачную модель ошибок:
upstream exception
↓
downstream может обработать его
Нельзя проектировать producer так, чтобы он скрыто продолжал emit после неуспешного downstream emission.
На собеседовании важно понимать направление:
source → operators → collector
и место каждого catch.
buffer меняет topology pipeline
Без buffer концептуально:
producer → operator → collector
выполняются как тесно связанный sequential pipeline.
С:
.buffer()
появляется дополнительная coroutine/channel boundary:
upstream coroutine
↓
buffer
↓
downstream coroutine
Поэтому buffer — не просто «список элементов в памяти».
flowOn тоже создаёт boundary
flow
.map(::A)
.flowOn(Dispatchers.IO)
.map(::B)
Нужно обеспечить передачу элементов между upstream, работающим в другом context, и downstream.
Поэтому внутренне flowOn для смены dispatcher требует разделения выполнения частей pipeline.
Latest = cancellation
Ключевая идея:
latest semantics
=
previous work cancellation
Например:
mapLatest { ... }
collectLatest { ... }
flatMapLatest { ... }
Это не просто «игнорирование старого результата».
Предыдущая suspend-работа отменяется.
Поэтому код внутри должен корректно поддерживать cooperative cancellation.
Hot-flow lifetime
При stateIn/shareIn важно спрашивать:
- Какой
scopeвладеет hot flow? - Когда upstream запускается?
- Когда останавливается?
- Что кешируется?
- Что получит новый subscriber?
Например:
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5_000),
initialValue = UiState(),
)
Lifetime hot flow ограничен viewModelScope, а lifecycle upstream дополнительно регулируется SharingStarted.
State vs Event
Это архитектурный вопрос, а не вопрос выбора класса.
State:
"Что истинно сейчас?"
Event:
"Что произошло?"
Пример состояния:
data class UiState(
val paymentStatus: PaymentStatus
)
Если payment completion критично восстановить после recreation, хранить его только как transient:
SharedFlow<PaymentCompleted>
может быть неверной моделью.
replay — не гарантия durable delivery
SharedFlow(replay = 1) хранит replay cache в памяти hot-flow instance.
Это не persistent event storage.
При уничтожении владельца потока данные исчезнут.
Для durable state нужны соответствующие persistence mechanisms.
StateFlow и одинаковые значения
state.value = User("Alex")
state.value = User("Alex")
Если значения equal, второй update не обязан породить новую emission.
Поэтому нельзя использовать StateFlow как механизм:
«Мне нужно обязательно доставить каждый одинаковый сигнал».
Это другая семантика.
Практика
Вопросы на собеседовании
Чем Flow отличается от StateFlow?
Ответ:
Обычный Flow чаще cold: каждый collector запускает upstream отдельно. StateFlow — hot state holder с обязательным initial value, синхронным .value и equality-based conflation.
StateFlow vs SharedFlow?
StateFlow специализирован под текущее состояние.
SharedFlow — общий hot broadcast stream с configurable replay/buffer semantics.
combine vs zip?
zip сопоставляет элементы попарно.
combine после первого значения каждого source пересчитывает результат при изменении любого source с использованием последних значений остальных.
flatMapLatest зачем?
Когда outer value определяет актуальный inner source.
selectedUserId
.flatMapLatest(repository::observeUser)
При новом ID предыдущий inner Flow отменяется.
collectLatest vs flatMapLatest?
collectLatest отменяет обработку предыдущего элемента внутри collector.
flatMapLatest отменяет предыдущий inner Flow.
Что делает flowOn?
Меняет coroutine context upstream относительно своего положения. Downstream context collector он не переключает.
Что произойдёт при двух collectors cold Flow?
Upstream запустится отдельно для каждого collector.
Когда нужен shareIn?
Когда несколько subscribers должны разделять один upstream вместо повторного запуска cold producer.
Почему StateFlow требует initial value?
Потому что он моделирует текущее состояние и должен всегда иметь доступное .value.
Почему update лучше read-modify-write через .value?
Потому что update предоставляет атомарную операцию изменения значения и не теряет конкурентные изменения так, как наивный read-modify-write.
Чем conflate отличается от collectLatest?
conflate может пропускать промежуточные элементы.
collectLatest отменяет обработку предыдущего элемента.
Почему catch может не поймать exception?
catch работает с upstream относительно своего положения. Exception, возникший downstream, например в последующем collect, находится вне его области.
Зачем нужен awaitClose в callbackFlow?
Для cleanup callback/listener/resource при закрытии или cancellation Flow.
Flow vs Channel?
Channel — coroutine communication primitive. Flow — stream abstraction с declarative operators и collection semantics.
Flow vs suspend function?
suspend fun — обычно одна асинхронная операция/результат.
Flow — последовательность значений во времени.
Что означает backpressure в Flow?
Медленный downstream способен естественно замедлять upstream через suspension emit.
Что меняет buffer?
Разделяет upstream/downstream дополнительной buffering/coroutine boundary, позволяя им работать более независимо.
Что означает hot Flow?
Его lifecycle и producer не определяются отдельным вызовом collect; один stream может обслуживать несколько subscribers.
Что получит новый subscriber SharedFlow(replay = 0)?
Он не получит старые emissions из replay cache. Он будет получать последующие значения в соответствии с конфигурацией потока.
Почему не каждый UI event должен быть SharedFlow?
Потому что часть «events» на самом деле является состоянием, которое UI должен уметь восстановить после временного отсутствия collector или recreation.
Шпаргалка
| Задача | Инструмент |
|---|---|
| Одно async значение | suspend fun |
| Поток значений | Flow |
| UI state | StateFlow |
| Общий hot broadcast stream | SharedFlow |
| Несколько states → один | combine |
| Попарно объединить | zip |
| Новый input отменяет transform | mapLatest |
| Новый source отменяет старый source | flatMapLatest |
| Новый value отменяет collector work | collectLatest |
| Задержать частые изменения | debounce |
| Удалить соседние дубли | distinctUntilChanged |
| Развязать producer/consumer | buffer |
| Нужны свежие значения | conflate |
| Изменить upstream dispatcher | flowOn |
| Cold → SharedFlow | shareIn |
| Cold → StateFlow | stateIn |
| Callback API → Flow | callbackFlow |
| Retry upstream | retry / retryWhen |
| Side effect на элементе | onEach |
| Collect в scope | launchIn |
Главная mental model
Для любого Flow-кода ответь на восемь вопросов:
- Кто производит значения?
- Cold это или hot?
- Когда producer запускается?
- Когда producer останавливается?
- Что происходит, если consumer медленнее producer?
- В каком CoroutineContext выполняется каждая часть pipeline?
- Что отменяется при новом значении?
- Что увидит новый subscriber?
Если ты уверенно отвечаешь на них, ты понимаешь Flow существенно глубже, чем на уровне запоминания операторов.
Финальная схема
COLD
│
Flow<T>
│
operators
│
┌──────────┴──────────┐
│ │
stateIn() shareIn()
│ │
▼ ▼
StateFlow<T> SharedFlow<T>
│ │
└──────────┬──────────┘
│
HOT
│
▼
Android UI
Практика перед собеседованием
Попробуй без запуска кода объяснить для каждого примера:
- когда стартует upstream;
- сколько раз он стартует;
- где происходит suspension;
- что произойдёт при cancellation;
- какой dispatcher используется;
- какие значения могут быть потеряны;
- что увидит новый subscriber;
- завершится ли Flow;
- какое исключение и где будет обработано.
Особенно полезно потренироваться на комбинациях:
flowOn + buffer
debounce + mapLatest
combine + stateIn
flatMapLatest + catch
SharedFlow + replay + extraBufferCapacity
StateFlow + update
Это тип задач, который хорошо отделяет знание API от понимания механики Kotlin Flow.