Explain the ordering and concurrency guarantees of flatMapConcat versus flatMapMerge, including the concurrency argument and its default.
answer
- Concat: order preserved, no overlap
- Merge: interleaved across inner flows
- DEFAULT_CONCURRENCY = 16
- concurrency arg caps active inner flows
- Within one inner flow order kept
basics
~20 sConcat processes inner flows one at a time, so output keeps the input order. Merge processes several at once, so faster inner flows can finish first and outputs interleave. Merge lets you cap how many run together; the default is 16.
solid answer
~40 sflatMapConcat subscribes to inner flow N+1 only after inner flow N completes. Emissions therefore appear strictly in upstream order — deterministic but with no parallelism, so total time is the sum of inner-flow durations. flatMapMerge subscribes to up to `concurrency` inner flows simultaneously and emits each value as soon as it's ready; the relative order of values from different inner flows is non-deterministic, though values within a single inner flow keep their internal order. The `concurrency` parameter defaults to `DEFAULT_CONCURRENCY` (16) and can be set per call: `flatMapMerge(concurrency = 4) { ... }`. Setting it to 1 makes merge behave like concat in throughput but it's clearer to use flatMapConcat. Unbounded concurrency risks resource exhaustion, which is why a finite default exists.
code
kotlin · 10 linesimport kotlinx.coroutines.flow.*
import kotlinx.coroutines.delay
fun req(id: Int, ms: Long) = flow { delay(ms); emit("r$id") }
suspend fun bounded() {
(1..100).asFlow()
.flatMapMerge(concurrency = 8) { req(it, 50L) }
.collect(::println) // at most 8 requests in flight at once
}go deeper
States concat is ordered/sequential and merge is concurrent/unordered.
Knows the concurrency parameter, its default of 16, and reasons about wall-clock time differences.
Explains how upstream is suspended at the concurrency cap and the resource trade-offs of raising it.
Weighs bounded concurrency against downstream/system limits and discusses tuning the global default vs per-call values.
## flatMapConcat ordering `flatMapConcat` is **sequential**: it collects each inner flow to completion before starting the next. Consequences: - Output order == input order (deterministic). - No overlap: wall-clock time ≈ sum of all inner-flow durations. - Backpressure is natural — only one inner flow is active. ## flatMapMerge ordering `flatMapMerge` is **concurrent**: it can collect several inner flows at the same time and forwards values as soon as they arrive. - Values from **different** inner flows interleave non-deterministically. - Values **within one** inner flow keep their own order. - Wall-clock time ≈ max duration among the concurrently running inner flows (subject to the concurrency cap). ## The concurrency argument ```kotlin public fun <T, R> Flow<T>.flatMapMerge( concurrency: Int = DEFAULT_CONCURRENCY, transform: suspend (value: T) -> Flow<R> ): Flow<R> ``` - `DEFAULT_CONCURRENCY` is **16**. - It can be overridden globally with the JVM system property `kotlinx.coroutines.flow.defaultConcurrency`. - Pass `concurrency = N` to limit how many inner flows run at once; upstream is suspended once N inner flows are active until one finishes. - `concurrency = 1` reduces merge to effectively sequential behaviour (but prefer `flatMapConcat` for clarity). ```kotlin import kotlinx.coroutines.flow.* import kotlinx.coroutines.delay fun job(id: Int, ms: Long) = flow { delay(ms); emit(id) } suspend fun run() { val src = flowOf(1 to 300L, 2 to 100L, 3 to 200L) src.flatMapConcat { (id, ms) -> job(id, ms) }.collect(::println) // 1,2,3 (ordered, ~600ms total) src.flatMapMerge { (id, ms) -> job(id, ms) }.collect(::println) // 2,3,1 (by speed, ~300ms total) } ``` ## When to choose which - Need ordered results / dependent steps → **flatMapConcat**. - Independent work you want to parallelise (bounded) → **flatMapMerge** with a sensible `concurrency`. ## Gotcha Never assume merge output order. Sort or tag results yourself if you need order back.
- What is the wall-clock time difference between concat and merge for 3 inner flows of 100ms each (concurrency >= 3)?Concat ~300ms (sum), merge ~100ms (max), assuming enough concurrency to run all three together.
- How do you change the merge default for the whole JVM?Set the system property kotlinx.coroutines.flow.defaultConcurrency.
saying these in an interview costs you the question
- Saying merge guarantees input order
- Not knowing the default concurrency is 16
- Claiming concat runs inner flows in parallel with concurrency=1
- Thinking concurrency=0 is valid (it must be >= 1)
- Believing within-one-inner-flow order can be reordered by merge