Lectures

Senior

Kotlin Flow: From Pipelines to Hot State

Flow execution, operators, hot state, lifecycle collection, delivery policies, and interview practice.

flowcoroutinesstateflowsharedflowandroid
On this page

Follow values from producer to UI and choose operators by execution, cancellation, and delivery semantics. All supplied examples remain read-only fragments; they were not compiled or executed during preparation. Imports, dependencies, and experimental opt-ins depend on the project's Kotlin and coroutine versions.

Learning objectives:

  • trace cold execution, context boundaries, and backpressure;
  • choose combination, flattening, buffering, and latest semantics;
  • define state, broadcast, sharing, and lifecycle ownership;
  • explain failures, cancellation, callback cleanup, and interview scenarios.

Theory

An asynchronous stream

Flow<T> represents zero or more asynchronous values over time. A suspend function usually supplies one result. A repository stream can feed ViewModel state and a Compose UI.

suspend fun getUser(): User
fun observeUsers(): Flow<List<User>>
time ──────────────────────────────>

       users1     users2     users3
Flow ───●──────────●──────────●─────>
Room / DataStore / Repository
              ↓
             Flow
              ↓
           ViewModel
              ↓
      StateFlow<UiState>
              ↓
collectAsStateWithLifecycle()
              ↓
           Compose

Producer, operators, collector

The producer emits, intermediate operators transform or select, and a terminal collector consumes. Constructing an operator chain does not itself collect it.

Producer
   │
   ▼
flow { emit(...) }
   │
   ▼
map
   │
   ▼
filter
   │
   ▼
debounce
   │
   ▼
Collector
repository.observeUsers()
    .map { users -> users.filter(User::isActive) }
    .distinctUntilChanged()
    .collect { users ->
        render(users)
    }

Cold execution

A regular cold flow starts its upstream for each collector. Creating it starts no producer. Two collectors normally perform upstream work twice.

val numbers = flow {
    println("Started")
    emit(1)
    emit(2)
    emit(3)
}
val numbers = flow { ... }
numbers.collect {
    println(it)
}
numbers.collect { println("A: $it") }
numbers.collect { println("B: $it") }
Collector A ──► producer execution #1
Collector B ──► producer execution #2

Official reference: Cold execution API.

Creating flows

Use flow for suspending production, flowOf for fixed values, and asFlow to adapt supported collections or sequences. A builder can invoke suspend functions.

fun numbers(): Flow<Int> = flow {
    emit(1)
    emit(2)
    emit(3)
}
fun users(): Flow<User> = flow {
    val users = api.loadUsers()

    users.forEach {
        emit(it)
    }
}
flowOf(1, 2, 3)

listOf(1, 2, 3).asFlow()

Terminal operations

collect suspends until completion. first and firstOrNull select the first value; single and singleOrNull enforce single-value semantics; toList, reduce, and fold aggregate. launchIn starts collection in an explicit scope. Two infinite collect calls in sequence never reach the second call: launch separate child collectors or combine the streams.

val result = repository.users()
    .map { ... }
    .filter { ... }
result.collect {
    println(it)
}
first()
firstOrNull()

single()
singleOrNull()

toList()

reduce()
fold()

launchIn(scope)
viewModelScope.launch {
    repository.users().collect {
        println(it)
    }

    println("DONE")
}
viewModelScope.launch {
    flowA.collect { ... }
    flowB.collect { ... }
}
viewModelScope.launch {
    flowA.collect { ... }
}

viewModelScope.launch {
    flowB.collect { ... }
}

Official reference: Terminal operations API.

Sequential execution

Ordinary flow execution is sequential: emit passes through downstream processing before proceeding. Flow does not imply automatic parallel work.

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)

Backpressure, buffering, conflation

Suspending emit coordinates producer and consumer. buffer introduces a coroutine and channel boundary with explicit capacity and overflow semantics. conflate discards intermediate pending values when a consumer is slow; choose it only when freshness matters more than every value.

flow {
    repeat(1000) {
        emit(it)
    }
}.collect {
    delay(1000)
}
Producer                Collector

emit(1) ──────────────► process 1
   │                       │
   │ suspend               │
   │                       │
emit(2) ◄──────────────────┘
flow
    .buffer()
    .collect {
        slowOperation(it)
    }
Producer        Buffer        Consumer

   ● ─────────► [●]
   ● ─────────► [● ●] ──────► slow
   ● ─────────► [● ●]
flow
    .conflate()
    .collect {
        render(it)
    }
Producer:   1 2 3 4 5 6 7 8
Collector:  1 ───── 4 ───── 8

Latest operators

collectLatest cancels the previous collector action when a new value arrives. mapLatest cancels the previous transformation. conflate skips intermediate pending values; latest operators cancel already-started work, which must cooperate with cancellation.

flow.collectLatest { value ->
    process(value)
}
"a"
 └── process("a") ─────X

"ab"
 └── process("ab") ────X

"abc"
 └── process("abc") ───── DONE
query
    .mapLatest { query ->
        repository.search(query)
    }

Official reference: Latest operators API.

Debounce and adjacent duplicates

debounce forwards after a quiet interval, for example a 300 ms search delay. distinctUntilChanged removes adjacent equivalent values: 1,1,1,2,2,3,3,2 becomes 1,2,3,2. Select intervals from product requirements.

query.debounce(300)
A──50ms──Al──80ms──Ale──100ms──Alex
                              │
                              │ 300 ms
                              ▼
                            emit
query
    .debounce(300)
    .filter { it.length >= 2 }
    .distinctUntilChanged()
    .mapLatest(repository::search)
1 1 1 2 2 3 3 2
1 2 3 2

Errors and cancellation

catch handles exceptions upstream of its position, not downstream collector failures. Coroutine cancellation is normal control flow. Broad exception handling must rethrow CancellationException rather than convert it to an error state.

flow {
    emit(api.loadUsers())
}
.catch {
    emit(emptyList())
}
flow
    .catch { ... }
    .collect {
        throw RuntimeException()
    }
upstream
   │
   ▼
 catch
   │
   ▼
collector
   │
   X exception
val job = scope.launch {
    flow.collect { ... }
}

job.cancel()
try {
    ...
} catch (e: CancellationException) {
    throw e
} catch (e: Throwable) {
    ...
}

Official reference: Errors and cancellation API.

Context and flowOn

flow preserves the collector context. Emitting from a different withContext context inside flow is deliberately invalid. flowOn changes upstream context only: transformations before it use that context, while later operations and the collector retain their context.

flow {
    withContext(Dispatchers.IO) {
        emit(loadData())
    }
}
flowOn()
repository.users()
    .map(::heavyTransformation)
    .flowOn(Dispatchers.Default)
    .collect {
        render(it)
    }
flow
    .map { A(it) }
    .flowOn(Dispatchers.IO)
    .map { B(it) }
    .collect { C(it) }
Dispatchers.IO:
    producer
       ↓
      A()

---------- flowOn boundary ----------

Main:
      B()
       ↓
      C()

Official reference: Context and flowOn API.

Combining and pairing

zip pairs corresponding emissions. combine waits for initial values from all sources, then uses their latest values when a source changes. Current UI state usually needs the latter semantics.

a.zip(b) { x, y ->
    x to y
}
A: A1 ─── A2 ─── A3
B: B1 ─────── B2 ─── B3

   A1B1      A2B2    A3B3
combine(
    userFlow,
    settingsFlow,
) { user, settings ->
    UiState(user, settings)
}
User:
 U1────────────U2────────>

Settings:
      S1──S2─────────────>

combine:
      U1S1
          U1S2
               U2S2

Official reference: Combining and pairing API.

Flattening streams

Mapping an ID to repository.observe produces nested flows. flatMapConcat collects inner flows sequentially; flatMapMerge allows concurrent inner flows and interleaving; flatMapLatest cancels the previous inner observation on a new ID. Selected users, accounts, chats, venues, and search queries often require latest selection.

val ids: Flow<UserId>

fun observeUser(id: UserId): Flow<User>
ids.map {
    observeUser(it)
}
Flow<Flow<User>>
A → [1 2 3]
            B → [4 5 6]
                       C → [7 8]
A → [────────]
B →   [────────]
C →     [────────]
selectedUserId
    .flatMapLatest(repository::observeUser)
User 1: [────────────X]
User 2:        [──────────X]
User 3:               [────────────>

Official reference: Flattening streams API.

Hot streams

A hot stream can exist independently of an individual collector. StateFlow and SharedFlow are hot; lifecycle and ownership still determine how their producers run.

               ┌── Collector A
Producer ──────┼── Collector B
               └── Collector C
StateFlow<T>
SharedFlow<T>

StateFlow and atomic updates

StateFlow requires an initial value, exposes current value, and gives a new collector the latest state. Equality conflation suppresses equivalent updates and can skip intermediate states. Expose a read-only view, use immutable state copies, and use update for atomic read-modify-write. Mutating a list held inside the same state object can hide changes.

private val _state =
    MutableStateFlow(UiState())

val state: StateFlow<UiState> =
    _state.asStateFlow()
_state.value = newState
_state.update {
    it.copy(isLoading = true)
}
state.value
current state = S3

new collector
     │
     ▼
     S3
data class UiState(
    val loading: Boolean = false,
    val users: List<User> = emptyList(),
)
StateFlow<MutableList<User>>
_state.update {
    it.copy(users = it.users + user)
}
_state.value = _state.value.copy(
    counter = _state.value.counter + 1
)
_state.update {
    it.copy(counter = it.counter + 1)
}

Official reference: StateFlow and atomic updates API.

SharedFlow and delivery policy

SharedFlow broadcasts values without a mandatory current value. replay controls cached values for new subscribers; extraBufferCapacity and onBufferOverflow control buffering with subscribers. Policies include SUSPEND, DROP_OLDEST, and DROP_LATEST. emit can suspend; tryEmit returns immediate acceptance status. Without subscribers, only replay values are retained; extra buffer capacity does not provide a durable offline queue. A successful tryEmit does not prove delivery to a consumer.

private val _events =
    MutableSharedFlow<Event>()

val events =
    _events.asSharedFlow()
_events.emit(Event.LoggedOut)
MutableSharedFlow<Event>(
    replay = 1
)
A → B → C

new subscriber:

replay = 0 → nothing
replay = 1 → C
replay = 2 → B, C
StateFlow:
"What is current?"

SharedFlow:
"Which values are broadcast?"
replay
extraBufferCapacity
onBufferOverflow
_events.emit(event)
val accepted = _events.tryEmit(event)
MutableSharedFlow<Event>(
    replay = 0,
    extraBufferCapacity = 64,
    onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
BufferOverflow.SUSPEND
BufferOverflow.DROP_OLDEST
BufferOverflow.DROP_LATEST

Sharing cold upstream

shareIn creates a scope-owned SharedFlow with replay; stateIn creates a StateFlow with an initial value. Eagerly starts immediately, Lazily on the first subscriber, and WhileSubscribed controls upstream by active subscriptions. A five-second stop timeout can bridge a brief UI recreation gap; choose it deliberately.

Collector A ──► producer A
Collector B ──► producer B
val shared = repository.users()
    .shareIn(
        scope = viewModelScope,
        started = SharingStarted.WhileSubscribed(),
        replay = 1,
    )
                  ┌── Collector A
producer ─ shared ┤
                  └── Collector B
val state: StateFlow<UiState> =
    repository.observeUsers()
        .map { users ->
            UiState(users = users)
        }
        .stateIn(
            scope = viewModelScope,
            started = SharingStarted.WhileSubscribed(5_000),
            initialValue = UiState(),
        )
SharingStarted.WhileSubscribed(5_000)

Official reference: Sharing cold upstream API.

Android and Compose collection

Combine repository data and selection into ViewModel state, then stateIn with an explicit scope and subscription policy. Compose uses collectAsStateWithLifecycle; Views can use repeatOnLifecycle at STARTED. Collection and upstream ownership are distinct decisions.

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(),
        )
}
val uiState by viewModel.uiState
    .collectAsStateWithLifecycle()
Repository
    ↓
Flow
    ↓
ViewModel
    ↓
StateFlow<UiState>
    ↓
collectAsStateWithLifecycle()
    ↓
Compose State
    ↓
Recomposition
lifecycleScope.launch {
    repeatOnLifecycle(Lifecycle.State.STARTED) {
        viewModel.state.collect {
            render(it)
        }
    }
}

Adapting callbacks

callbackFlow bridges callback APIs into Flow. Register the callback and use awaitClose to unregister it on cancellation or closure. Missing cleanup leaks resources or duplicates callbacks. Inspect trySend results when delivery requirements demand handling a closed or full channel.

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
collector cancelled
       ↓
callbackFlow closes
       ↓
awaitClose
       ↓
unregister(listener)

Official reference: Adapting callbacks API.

Flow, Channel, Sequence, suspend

Channel is a communication primitive with send/receive and queue or rendezvous semantics. Flow is a stream abstraction with transformations and collection; channelFlow, callbackFlow, buffer, and flowOn may use channels internally. Sequence is synchronous pull-based processing. Prefer a suspend function for one asynchronous result when ongoing changes are unnecessary.

Flow
→ stream abstraction / declarative pipeline

Channel
→ coroutine communication primitive
Sequence → synchronous
Flow     → asynchronous + suspend
flow {
    delay(1000)
    emit(1)
}
suspend fun login(): User
fun observeUser(): Flow<User>

Retrying upstream

retry restarts upstream work; retryWhen can condition retries and backoff. Repeating non-idempotent side effects can duplicate actions, so classify retryable failures and operations before enabling retries.

repository.load()
    .retry(3) {
        it is IOException
    }
flow {
    chargeCreditCard()
    emit(...)
}
.retryWhen { cause, attempt ->
    if (
        cause is IOException &&
        attempt < 3
    ) {
        delay(1000L * (attempt + 1))
        true
    } else {
        false
    }
}

Official reference: Retrying upstream API.

Architecture mistakes

Avoid sequential infinite collectors, one-shot operations wrapped in Flow without a need, mutable state hidden inside StateFlow, automatically selecting SharedFlow for every UI action, duplicate expensive cold upstream collections, assuming flowOn changes downstream, and swallowing cancellation.

flowA.collect { ... }
flowB.collect { ... }
fun login(): Flow<User>
suspend fun login(): User
StateFlow<MutableList<User>>
expensiveFlow.collect(...)
expensiveFlow.collect(...)
shareIn(...)
stateIn(...)

Senior mechanics

Suspending emit

Suspension supplies backpressure and preserves ordinary sequential processing.

flow {
    emit(value)
}
upstream and downstream
coordinate through suspension

Context preservation

Do not emit from an unrelated context in flow. Use flowOn for upstream context, or channelFlow for concurrent producers.

flowOn(Dispatchers.IO)

Exception transparency

Do not continue emitting after a downstream failure. catch has a directional boundary and must not hide collector failures.

upstream exception
      ↓
downstream can handle it
source → operators → collector

Buffer changes topology

A buffer creates a coroutine/channel boundary, not merely a list of saved values.

producer → operator → collector
.buffer()
upstream coroutine
       ↓
     buffer
       ↓
downstream coroutine

flowOn boundaries

Changing dispatchers with flowOn introduces a coroutine boundary and transfer between contexts. A flowOn call that does not change dispatcher need not create another coroutine.

flow
    .map(::A)
    .flowOn(Dispatchers.IO)
    .map(::B)

Latest means cancellation

Latest operators cancel previous work. CPU loops and external integrations must cooperate with cancellation; hiding a result alone is different.

latest semantics
      =
previous work cancellation
mapLatest { ... }
collectLatest { ... }
flatMapLatest { ... }

Hot stream lifetime

Separate the owning scope, upstream start/stop policy, replay cache, and individual subscriptions. viewModelScope ownership does not make every upstream run forever.

.stateIn(
    scope = viewModelScope,
    started = SharingStarted.WhileSubscribed(5_000),
    initialValue = UiState(),
)

State and events

Current state answers what is true now; an event describes what occurred. Restorable success state and a transient navigation instruction have different delivery requirements.

"What is true now?"
"What happened?"
data class UiState(
    val paymentStatus: PaymentStatus
)
SharedFlow<PaymentCompleted>

Replay is not durable storage

Replay is an in-memory cache owned by the stream. Owner destruction loses it; durable delivery or recovery requires an explicit persistence design.

Equivalent StateFlow values

Two equal Loading states do not guarantee two emissions. StateFlow models current information, not delivery of every equal signal.

state.value = User("Alex")
state.value = User("Alex")

Practice

Explain the expected behavior of every supplied example before experimenting in a configured project.

Interview questions

Flow versus StateFlow

Flow is a general stream; StateFlow is a hot current-state stream with an initial value and equality conflation.

StateFlow versus SharedFlow

StateFlow specializes current state. SharedFlow provides configurable broadcast replay and buffering without mandatory initial/current value semantics.

combine versus zip

combine uses latest values after each source emits; zip pairs corresponding values.

Why flatMapLatest?

Cancel an old inner observation when the selected ID or query changes.

selectedUserId
    .flatMapLatest(repository::observeUser)

collectLatest versus flatMapLatest

The first cancels the previous collector action; the second switches and cancels the previous inner flow.

What does flowOn change?

Upstream context only, not downstream collector context.

Two collectors of a cold flow

Normally two independent upstream executions.

When to shareIn?

Share one upstream among subscribers, with an explicit scope, start policy, and replay requirement.

Why an initial StateFlow value?

A current state must exist before any upstream emission.

Why atomic update?

It avoids races in a separate read-modify-write through value.

conflate versus collectLatest

Conflate skips pending intermediate values; collectLatest cancels processing of an earlier value.

Why might catch miss an exception?

It handles upstream failures, not downstream collector failures.

Why awaitClose?

Unregister callbacks and release resources when the callback flow closes or collection cancels.

Flow versus Channel

A stream abstraction versus a communication primitive; some Flow operators use channels internally.

Flow versus a suspend function

Ongoing values versus one asynchronous result.

What is backpressure?

A slow consumer can suspend production through emit.

What does buffer change?

It decouples producer and consumer through a coroutine/channel boundary and a delivery policy.

What is hot Flow?

Its existence is independent of an individual collector; ownership still governs production.

New subscriber with replay zero

No cached past values; future delivery follows subscription and buffering semantics.

Why not SharedFlow for every UI action?

Define absence, loss, replay, restoration, and durable delivery requirements first.

Operator cheat sheet

Requirement Tool
One asynchronous result suspend function
Changing values Flow
Current UI state StateFlow
Broadcast values SharedFlow
Combine current states combine
Pair corresponding values zip
Cancel prior transformation mapLatest
Switch inner observation flatMapLatest
Cancel collector action collectLatest
Wait for a quiet interval debounce
Suppress adjacent duplicates distinctUntilChanged
Decouple producer and consumer buffer
Keep fresh pending values conflate
Change upstream context flowOn
Share a cold flow shareIn
Expose shared current state stateIn
Adapt callbacks callbackFlow
Retry upstream retry / retryWhen
Run an action per value onEach
Collect in a scope launchIn

Ask who produces, whether the stream is cold or hot, when it starts and stops, what a slow consumer does, where each stage runs, what a new value cancels, and what a new subscriber receives.

                    COLD
                     │
                  Flow<T>
                     │
                  operators
                     │
          ┌──────────┴──────────┐
          │                     │
       stateIn()             shareIn()
          │                     │
          ▼                     ▼
     StateFlow<T>         SharedFlow<T>
          │                     │
          └──────────┬──────────┘
                     │
                    HOT
                     │
                     ▼
                 Android UI

Work through pipeline combinations

For each supplied pipeline, identify start conditions, collector count, suspension, cancellation, dispatchers, dropped values, replay, completion, and error boundaries. The expected result is an explanation of semantics; runtime output has not been verified.

flowOn + buffer
debounce + mapLatest
combine + stateIn
flatMapLatest + catch
SharedFlow + replay + extraBufferCapacity
StateFlow + update

Further reading