Senior
Kotlin Flow: основы
Холодные и горячие потоки Kotlin Flow, context, операторы, backpressure, sharing и UI state.
На этой странице
Это русскоязычная подготовленная версия предоставленного конспекта. Все исходные кодовые блоки сохранены как read-only примеры. Примеры не компилировались и не запускались в ходе подготовки.
Теория
Зачем Flow, если есть suspend
suspend:
один запрос → один результат
suspend fun getUser(): User
Flow:
0..N значений во времени
fun observeBookings(): Flow<List<Booking>>
Официальная документация: Документация Kotlin Coroutines.
Flow — асинхронный stream
val numbers: Flow<Int> = flow {
emit(1)
emit(2)
emit(3)
}
numbers.collect { value ->
println(value)
}
emit() производит значение, collect() потребляет.
Официальная документация: Flow API.
Suspend-based API
Концептуально:
interface Flow<out T> {
suspend fun collect(
collector: FlowCollector<T>
)
}
interface FlowCollector<in T> {
suspend fun emit(value: T)
}
Поэтому Flow естественно интегрирован с cancellation и CoroutineContext.
Официальная документация: Документация Kotlin Coroutines.
Cold Flow
val flow = flow {
println("START")
emit(1)
emit(2)
}
Пока никто не делает collect, producer не запускается.
Каждый новый collector обычного cold Flow запускает upstream заново.
Официальная документация: Flow API.
Flow vs Sequence
Sequence<T> → lazy synchronous stream
Flow<T> → lazy asynchronous suspend-aware stream
В Flow можно использовать suspend operations.
Официальная документация: Документация Kotlin Coroutines.
Context collector'а
Без специальных операторов upstream выполняется в coroutine context collector'а.
Отсюда необходимость flowOn.
flowOn
flow
.flowOn(Dispatchers.Default)
.collect {
render(it)
}
flowOn влияет на upstream.
flow
.map { operationA(it) }
.flowOn(Dispatchers.Default)
.map { operationB(it) }
.collect { ... }
Default:
flow
↓
map A
↓
flowOn
──────────────
collector context:
map B
↓
collect
Официальная документация: Dispatchers API.
Context preservation
Не следует произвольно делать:
flow {
withContext(Dispatchers.IO) {
emit(load())
}
}
Для upstream context используется:
flow {
emit(load())
}.flowOn(Dispatchers.IO)
Официальная документация: Dispatchers API.
Базовые operators
flow
.map { it.toUiModel() }
.filter { it.isActive }
.collect { render(it) }
transform позволяет создать 0..N outputs на один input.
onStart
flow
.map<Data, UiState> {
UiState.Content(it)
}
.onStart {
emit(UiState.Loading)
}
catch
flow
.catch { throwable ->
emit(UiState.Error(throwable))
}
.collect { state ->
render(state)
}
catch ловит upstream exceptions относительно своей позиции в chain.
Cancellation
Flow построен на coroutines:
scope cancelled
↓
collect cancelled
↓
upstream cancelled
Официальная документация: Flow API.
Backpressure
Flow по умолчанию может suspend'ить producer, пока consumer обрабатывает значение.
emit(1)
↓
collector processes 1
↓
emit(2)
Официальная документация: Документация Kotlin Coroutines.
buffer
flow
.buffer()
.collect {
slowOperation(it)
}
Producer и consumer получают возможность работать более независимо.
conflate
Если нужны только более свежие значения, промежуточные можно пропускать.
flow
.conflate()
.collect {
render(it)
}
Хороший пример — progress.
collectLatest
flow.collectLatest { value ->
process(value)
}
Новое значение отменяет обработку предыдущего.
A → process █████
X
B → process █████
mapLatest
queryFlow
.mapLatest { query ->
repository.search(query)
}
При новом query предыдущая transformation отменяется.
Официальная документация: Flow API.
debounce
query
.debounce(300)
.distinctUntilChanged()
.mapLatest { query ->
repository.search(query)
}
Подходит для search input, когда не нужно отправлять запрос на каждый символ.
distinctUntilChanged
A A A B B C
↓
A B C
StateFlow уже имеет equality-based conflation semantics.
Официальная документация: Flow API.
combine
combine(
user,
bookings,
) { user, bookings ->
ScreenState(user, bookings)
}
Использует latest значения sources.
User:
U1-------------U2---------
Bookings:
B1----B2---------B3---
combine:
U1B1--U1B2--U2B2--U2B3
zip
Сопоставляет emissions попарно.
A:
A1----A2----A3
B:
---B1----------B2----B3
zip:
---A1B1--------A2B2--A3B3
Ментальная модель:
combine = latest state
zip = corresponding pairs
flatMapLatest
userId.flatMapLatest { id ->
repository.observeUser(id)
}
При новом ID старая inner collection отменяется.
flatMapConcat
Inner flows выполняются последовательно.
A → █████
B → █████
C → █████
flatMapMerge
Inner flows могут выполняться concurrently.
A → █████████
B → █████
C → █████████
Итого:
flatMapConcat → sequential
flatMapMerge → concurrent
flatMapLatest → cancel previous
Cold vs Hot
Cold:
Collector A → Producer A
Collector B → Producer B
Hot:
Shared Producer
/ \
Collector A Collector B
flow {} обычно cold.
StateFlow и SharedFlow — hot.
Официальная документация: Flow API.
StateFlow
private val _state =
MutableStateFlow(UiState())
val state: StateFlow<UiState> =
_state.asStateFlow()
StateFlow:
- всегда имеет current value;
- требует initial value;
- новый collector сразу получает current value;
- имеет equality-based conflation;
- хорошо подходит для state.
Официальная документация: Flow API.
Atomic update
_state.update {
it.copy(isLoading = true)
}
Предпочтительно для read-modify-write обновлений при возможной concurrency.
SharedFlow
val events =
MutableSharedFlow<Event>()
Более общий hot broadcast stream без обязательной semantics «текущего состояния».
Официальная документация: Flow API.
replay
MutableSharedFlow<Event>(
replay = 2
)
Новый subscriber получает последние два replayed значения.
Но SharedFlow(replay = 1) всё равно не полностью равен StateFlow.
Официальная документация: Flow API.
State vs Event
State:
isLoading
user
bookings
selectedDate
Часто → StateFlow.
Event:
ShowSnackbar
NavigateBack
OpenUrl
Может быть SharedFlow/Channel, но сначала нужно определить delivery semantics.
Официальная документация: Flow API.
stateIn
val bookings =
repository.observeBookings()
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5_000),
initialValue = emptyList(),
)
Cold Flow
↓
stateIn
↓
StateFlow
Официальная документация: Flow API.
shareIn
val shared =
repository.events()
.shareIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(),
replay = 0,
)
Cold Flow
↓
shareIn
↓
SharedFlow
Официальная документация: Flow API.
Зачем sharing
Без sharing два collectors cold Flow могут запустить дорогой upstream дважды.
После sharing:
one upstream
│
┌───────┴───────┐
▼ ▼
Collector A Collector B
Официальная документация: Flow API.
SharingStarted
Частые варианты:
SharingStarted.Eagerly
SharingStarted.Lazily
SharingStarted.WhileSubscribed(...)
Для UI часто используется WhileSubscribed(timeout).
Timeout помогает не перезапускать upstream мгновенно при коротком configuration-change gap.
ViewModel pipeline
data class VenueUiState(
val venue: Venue? = null,
val bookings: List<Booking> = emptyList(),
val selectedDate: LocalDate = LocalDate.now(),
)
private val selectedDate =
MutableStateFlow(LocalDate.now())
val state: StateFlow<VenueUiState> =
combine(
repository.observeVenue(),
repository.observeBookings(),
selectedDate,
) { venue, bookings, date ->
VenueUiState(
venue = venue,
bookings = bookings,
selectedDate = date,
)
}.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5_000),
initialValue = VenueUiState(),
)
Официальная документация: Flow API.
Ментальная карта Flow
Flow
│
┌─────────┴─────────┐
│ │
Cold Hot
│ │
flow { } ┌──────┴──────┐
│ │
StateFlow SharedFlow
│ │
State broadcast
Cold Flow ──stateIn──► StateFlow
Cold Flow ──shareIn──► SharedFlow
Официальная документация: Flow API.
Практика
Вопросы для разбора
Для каждого примера нарисуйте иерархию Job, определите CoroutineContext и Dispatcher, отметьте точки suspension/cancellation, затем проследите failure вверх и cancellation вниз. Для Flow дополнительно определите cold/hot semantics, владельца collection, поведение медленного consumer, replay и последствия lifecycle stop/start.
Проверка архитектурного решения
Спросите: кто владеет этой работой, кто может её отменить, где обрабатывается ошибка, нужна ли каждая emission или только latest, и переживает ли результат пересоздание UI. Coroutine не переживает process death; для гарантированной фоновой работы нужен durable scheduler, например WorkManager.