Flow

Flow: Cold Streams of Values

A suspend fun returns one value; a Flow<T> emits many over time, such as Room 234 rows each time a table changes. A flow is cold, like an async function* generator: building it and chaining operators runs nothing, and each terminal operator (collect, first, reduce) runs the producer again from the start.

flow.kt: a cold flow runs once per terminal operatorKotlin
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun prices() = flow { log("flow started"); for (b in catalog) { delay(50); emit(b.price) } }
fun main() = runBlocking {
  val cheap = prices().filter { it < 20.0 }.map { "$%.2f".format(it) }
  log("flow built, nothing emitted yet")
  cheap.take(2).collect { log("got $it") }
  log("total %.2f".format(prices().flowOn(Dispatchers.Default).reduce { a, b -> a + b }))
}
Output
   0 ms [main @coroutine#1] flow built, nothing emitted yet
  36 ms [main @coroutine#1] flow started
 137 ms [main @coroutine#1] got $14.99
 299 ms [main @coroutine#1] got $18.75
 344 ms [DefaultDispatcher-worker-1 @coroutine#2] flow started
 666 ms [main @coroutine#1] total 134.74

take(2) cancelled the producer after the fourth book, and flowOn moved the upstream to a worker while reduce stayed on main. Use flowOn, never withContext, to move a producer: emitting from another context throws IllegalStateException.