Senior
Flow Fundamentals
Cold and hot Kotlin Flow, execution context, operators, backpressure, sharing, and ViewModel state pipelines.
On this page
Choose stream operators by ownership and delivery semantics. All supplied examples remain read-only coroutine or Android fragments; none was compiled or executed during preparation. Imports, dependencies, and opt-in requirements depend on the coroutine version used in an actual project.
Learning objectives:
- distinguish one asynchronous result from a stream;
- trace cold execution, context boundaries, and cancellation;
- choose buffering, latest, combination, and flattening semantics;
- model shared state and event delivery explicitly.
Theory
Flow or a suspend function
A suspend function represents one asynchronous result. Flow represents zero or more values over time, such as changing bookings. Choose a stream when ongoing changes matter, not merely because an API is asynchronous.
one request → one result
suspend fun getUser(): User
0..N values over time
fun observeBookings(): Flow<List<Booking>>
Official reference: Flow or a suspend function API.
An asynchronous stream
The producer emits values and the collector consumes them. Flow is asynchronous and suspend-aware; it is not automatically parallel.
val numbers: Flow<Int> = flow {
emit(1)
emit(2)
emit(3)
}
numbers.collect { value ->
println(value)
}
Official reference: An asynchronous stream API.
A suspend-based API
The conceptual Flow and FlowCollector interfaces use suspend collect and emit operations. This integrates stream processing with coroutine cancellation and context. The simplified interfaces below illustrate the contract rather than providing a full public implementation.
interface Flow<out T> {
suspend fun collect(
collector: FlowCollector<T>
)
}
interface FlowCollector<in T> {
suspend fun emit(value: T)
}
Official reference: A suspend-based API API.
Cold execution
Constructing a flow builder does not start its producer. Each new collector of a regular cold Flow normally starts its upstream again. Trace the START message when considering repeated collection.
val flow = flow {
println("START")
emit(1)
emit(2)
}
Flow and Sequence
Sequence is a lazy synchronous stream; Flow is a lazy asynchronous suspend-aware stream. Flow can use suspending production and collection operations.
Sequence<T> → lazy synchronous stream
Flow<T> → lazy asynchronous suspend-aware stream
The collector context
Without operators that introduce a context boundary, ordinary upstream execution uses the collector coroutine context. A slow or blocking upstream therefore needs deliberate scheduling, not an assumption that Flow runs in the background.
flowOn
flowOn affects upstream operators before its position, not downstream transformations or the collector. In the supplied chain, operationA uses the chosen upstream context, while operationB and collect remain in the collector context.
flow
.flowOn(Dispatchers.Default)
.collect {
render(it)
}
flow
.map { operationA(it) }
.flowOn(Dispatchers.Default)
.map { operationB(it) }
.collect { ... }
Default:
flow
↓
map A
↓
flowOn
──────────────
collector context:
map B
↓
collect
Official reference: flowOn API.
Context preservation
Emitting from an arbitrary different context inside flow violates its context-preservation invariant. The first snippet is deliberately incorrect. Move the upstream execution with flowOn instead of wrapping emit in withContext.
flow {
withContext(Dispatchers.IO) {
emit(load())
}
}
flow {
emit(load())
}.flowOn(Dispatchers.IO)
Basic operators
map transforms an input, filter selects values, and transform can emit zero to many outputs for each input. Intermediate operators construct a pipeline; a terminal operation activates collection.
flow
.map { it.toUiModel() }
.filter { it.isActive }
.collect { render(it) }
Official reference: Basic operators API.
onStart
onStart runs before upstream collection and can emit an initial Loading state. Its position and type must fit the rest of the pipeline.
flow
.map<Data, UiState> {
UiState.Content(it)
}
.onStart {
emit(UiState.Loading)
}
Official reference: onStart API.
catch
catch handles upstream exceptions relative to its position in the chain. It does not catch a failure in downstream collector code, and cancellation must not become a business failure.
flow
.catch { throwable ->
emit(UiState.Error(throwable))
}
.collect { state ->
render(state)
}
Official reference: catch API.
Flow cancellation
Canceling the owning scope cancels collection and its upstream work. Flow participates in coroutine ownership rather than creating an independent immortal producer.
scope cancelled
↓
collect cancelled
↓
upstream cancelled
Sequential backpressure
By default, processing an emitted value can suspend its producer until the consumer finishes. The pipeline naturally coordinates producer and consumer instead of accumulating every value without a bound.
emit(1)
↓
collector processes 1
↓
emit(2)
buffer
buffer introduces a buffering and coroutine boundary, allowing producer and consumer to progress more independently. Capacity and overflow policy must express delivery requirements.
flow
.buffer()
.collect {
slowOperation(it)
}
Official reference: buffer API.
conflate
If intermediate values are disposable, conflate lets a slow collector receive newer available values instead of every intermediate value. Progress display is one example; do not use this when every event matters.
flow
.conflate()
.collect {
render(it)
}
Official reference: conflate API.
collectLatest
A newer emission cancels the collector action processing the prior value. This cancels already-started work rather than merely hiding its result, and that work must cooperate with cancellation.
flow.collectLatest { value ->
process(value)
}
A → process █████
X
B → process █████
Official reference: collectLatest API.
mapLatest
A new query cancels the preceding transformation, making latest-result search semantics explicit. Cancellation of a coroutine does not prove that an external request never reached its server.
queryFlow
.mapLatest { query ->
repository.search(query)
}
Official reference: mapLatest API.
debounce
Debounce waits for a quiet period before forwarding a value. Combined with adjacent-duplicate suppression and mapLatest, it avoids searching for every keystroke. The 300 ms interval is an example product decision.
query
.debounce(300)
.distinctUntilChanged()
.mapLatest { query ->
repository.search(query)
}
Official reference: debounce API.
distinctUntilChanged
Suppress only consecutive equivalent values. StateFlow already has equality-based conflation, so applying the same operator to StateFlow does not add a new guarantee.
A A A B B C
↓
A B C
Official reference: distinctUntilChanged API.
combine
After each source has supplied a value, combine uses the latest values and can produce a new result when either changes. It fits coherent current UI state derived from several streams.
combine(
user,
bookings,
) { user, bookings ->
ScreenState(user, bookings)
}
User:
U1-------------U2---------
Bookings:
B1----B2---------B3---
combine:
U1B1--U1B2--U2B2--U2B3
Official reference: combine API.
zip
Zip matches corresponding emissions one-to-one, instead of combining each change with the other source’s latest state. Select it when pairing is the requirement.
A:
A1----A2----A3
B:
---B1----------B2----B3
zip:
---A1B1--------A2B2--A3B3
combine = latest state
zip = corresponding pairs
Official reference: zip API.
flatMapLatest
A new outer value cancels collection of the previous inner flow. This fits observing the currently selected user ID while abandoning the old observation.
userId.flatMapLatest { id ->
repository.observeUser(id)
}
Official reference: flatMapLatest API.
flatMapConcat
Inner flows are collected sequentially. A later inner flow waits for the previous one to complete, preserving sequential processing semantics.
A → █████
B → █████
C → █████
Official reference: flatMapConcat API.
flatMapMerge
Several inner flows can be collected concurrently, so emissions may interleave. It expresses a different delivery policy from sequential concatenation or canceling the prior flow.
A → █████████
B → █████
C → █████████
flatMapConcat → sequential
flatMapMerge → concurrent
flatMapLatest → cancel previous
Official reference: flatMapMerge API.
Cold and hot streams
A cold flow normally starts upstream independently for each collector. A hot stream can exist independently of an individual collector. StateFlow and SharedFlow are hot; flow builders are normally cold.
Collector A → Producer A
Collector B → Producer B
Shared Producer
/ \
Collector A Collector B
StateFlow
StateFlow always has a current value and requires an initial value. New collectors receive the current state; equality-based conflation can suppress equivalent updates and intermediate states. Expose a read-only StateFlow from its owner.
private val _state =
MutableStateFlow(UiState())
val state: StateFlow<UiState> =
_state.asStateFlow()
Official reference: StateFlow API.
Atomic state updates
update provides an atomic read-modify-write operation for concurrent changes. Prefer it to reading value and writing a copy as separate steps when updates can race.
_state.update {
it.copy(isLoading = true)
}
Official reference: Atomic state updates API.
SharedFlow
SharedFlow is a general hot broadcast stream without the mandatory current-value semantics of StateFlow. Decide replay and buffering from requirements rather than treating it as a drop-in event solution.
val events =
MutableSharedFlow<Event>()
Official reference: SharedFlow API.
Replay
A new subscriber receives the configured number of replayed values. SharedFlow with replay one still differs from StateFlow in current-value, initial-value, and equality behavior.
MutableSharedFlow<Event>(
replay = 2
)
Official reference: Replay API.
State and events
Loading, user, bookings, and selected date describe current state. Snackbar, navigation, and URL opening are actions with delivery requirements. Before choosing SharedFlow or Channel, decide what happens without a subscriber and whether loss or replay is acceptable.
isLoading
user
bookings
selectedDate
ShowSnackbar
NavigateBack
OpenUrl
stateIn
stateIn shares a cold upstream as a StateFlow in an explicit scope, with a started policy and initial value. The owner and sharing policy define upstream lifetime.
val bookings =
repository.observeBookings()
.stateIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(5_000),
initialValue = emptyList(),
)
Cold Flow
↓
stateIn
↓
StateFlow
Official reference: stateIn API.
shareIn
shareIn shares upstream as a SharedFlow, with an explicit scope, start policy, and replay count. Multiple downstream subscribers can use one upstream execution.
val shared =
repository.events()
.shareIn(
scope = viewModelScope,
started = SharingStarted.WhileSubscribed(),
replay = 0,
)
Cold Flow
↓
shareIn
↓
SharedFlow
Official reference: shareIn API.
Why sharing matters
Two collectors of expensive cold work can otherwise start it twice. Sharing can provide one upstream and fan its emissions out to several collectors. It also introduces a lifetime owner and policy that must be reviewed.
one upstream
│
┌───────┴───────┐
▼ ▼
Collector A Collector B
SharingStarted
Eagerly starts immediately; Lazily starts on the first subscriber; WhileSubscribed controls work while subscribers exist. A stop timeout can bridge a short configuration-change gap instead of restarting upstream immediately.
SharingStarted.Eagerly
SharingStarted.Lazily
SharingStarted.WhileSubscribed(...)
Official reference: SharingStarted API.
Practice
Trace the supplied pipeline and maps. Predict ownership, context, cancellation, and new-subscriber behavior before experimenting.
A ViewModel state pipeline
Purpose: derive one coherent VenueUiState from repository venue, bookings, and selected date. combine uses their latest values; stateIn gives the ViewModel-owned pipeline a current state and a subscription policy. Identify the initial state and every lifecycle owner in the supplied example.
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(),
)
The Flow mental map
Purpose: explain each branch of the supplied map. Distinguish cold and hot streams, state and broadcast, stateIn and shareIn. For every stream ask who starts upstream, who owns collection, how slow consumers are treated, and what a new subscriber sees.
Flow
│
┌─────────┴─────────┐
│ │
Cold Hot
│ │
flow { } ┌──────┴──────┐
│ │
StateFlow SharedFlow
│ │
State broadcast
Cold Flow ──stateIn──► StateFlow
Cold Flow ──shareIn──► SharedFlow
Choose a delivery policy
For a search stream, explain debounce, adjacent duplicates, and latest cancellation. For two current state streams, compare combine with pairwise zip. For a slow consumer, decide whether every value matters before choosing collect, buffer, conflate, or collectLatest. The expected result is a justified semantic choice; this page does not claim runtime verification.