skip to content

Where in a Flow chain does conflate() take effect, and what determines which dispatcher/coroutine the dropped-or-kept boundary sits on?

level: seniorimportance: nice to knowfreq 22%

answer

  1. conflate() = coroutine boundary at its position
  2. Upstream runs for every emission (even dropped)
  3. Downstream runs only on survivors
  4. Put conflate() before the expensive operator
  5. Dispatcher set by flowOn / collecting context, not conflate

basics

~20 s

conflate() splits the chain at the point you call it: everything before it runs in one coroutine, everything after in another, and dropping happens at that split. Operators upstream of conflate keep running; only values crossing the split get conflated.

solid answer

~50 s

Like buffer(), conflate() introduces a coroutine boundary at its position in the chain. Upstream operators (those before conflate()) run in the producer coroutine and keep emitting without waiting; downstream operators and the collector run in the consumer coroutine, and conflation drops intermediate values at the boundary. So placement matters: any map/filter you put before conflate() is evaluated on every emission (and its results may be dropped), whereas operators after conflate() only run for values that survive. The dispatcher of each side is governed by flowOn() and the collecting context — conflate() itself doesn't change dispatchers, it just inserts the conflated channel. Multiple conflate()/buffer() calls each add a boundary. A common optimization is to place conflate() right before an expensive downstream operator so the expensive work runs only on surviving values; conversely, putting cheap transforms upstream is fine even though they run on dropped values.

code

kotlin · 6 lines
kotlin
sensorReadings()
    .flowOn(Dispatchers.IO)    // producer dispatcher
    .conflate()                 // boundary: drop stale here
    .collect { reading ->       // collector context (e.g. main)
        renderLatest(reading)   // runs only for surviving values
    }

go deeper

for a junior

Understands conflate() sits at a point in the chain and separates producer from collector.

for a middle

Knows operators before conflate run for every emission while operators after run only on survivors.

for a senior

Optimizes placement (conflate before expensive ops) and explains that flowOn/collecting context, not conflate, set dispatchers.

for a principal

Reasons about channel fusion, end-to-end context design, and the CPU/latency trade-offs of operator ordering across the pipeline.

## conflate() is a fusion/boundary point `conflate()` inserts a **conflated channel** between its upstream and downstream, creating two coroutines that run **concurrently**: - **Upstream side** (operators *before* conflate): the producer. It keeps emitting; values pile into the 1-slot conflated channel, older ones dropped. - **Downstream side** (operators *after* conflate + the collector): the consumer. It reads surviving values from the channel. The drop happens **at the channel**, i.e. exactly where you wrote `conflate()`. ## Placement changes what runs on dropped values ```kotlin upstream .map { cheapTransform(it) } // runs for EVERY emission (some results dropped) .conflate() .map { expensiveRender(it) } // runs ONLY for surviving values .collect { ... } ``` Operators **before** `conflate()` execute for every upstream emission, even ones that will be dropped. Operators **after** `conflate()` execute only for values that made it through. So to avoid wasting CPU on doomed values, place `conflate()` **before** the expensive operator. ## Dispatchers: conflate() does not change them `conflate()` does **not** switch dispatchers by itself. The execution context is determined by: - **`flowOn(dispatcher)`** — changes the context of the **upstream** of the `flowOn` call. Combined with conflate's boundary, you control which dispatcher the producer runs on. - **the collecting coroutine's context** — governs the downstream/collector side. So the "dropped-or-kept boundary" sits between the upstream coroutine (possibly relocated by `flowOn`) and the collecting coroutine; conflate() just supplies the channel that joins them. ## Channel fusion Adjacent buffering operators are **fused**. For example `flowOn` already introduces a buffer; stacking `conflate()` adjacent to it adjusts the buffering rather than creating a redundant extra channel. Each explicit `conflate()`/`buffer()` you write, however, is a meaningful boundary you can reason about independently. ## Practical recipe ```kotlin sensorReadings() .flowOn(Dispatchers.IO) // produce off the main thread .conflate() // drop stale readings at the boundary .collect { reading -> // runs on the collecting (e.g. main) context renderLatest(reading) // only the freshest reading reaches here } ``` Here production happens on IO, conflation drops stale readings, and only the latest reading reaches the (main-thread) renderer.

  • If you put an expensive map() before conflate(), is its work wasted on dropped values?
    Yes — operators upstream of conflate() run for every emission, including values that conflate later drops. Move conflate() ahead of the expensive operator so it only runs on survivors.
  • Does conflate() change which dispatcher the producer runs on?
    No. conflate() only inserts the conflated channel boundary. The dispatcher is controlled by flowOn() for the upstream and by the collecting coroutine's context for the downstream.

saying these in an interview costs you the question

  • Thinks conflate() switches dispatchers like flowOn
  • Believes upstream operators skip dropped values
  • Says placement of conflate() in the chain doesn't matter
  • Cannot explain that drop happens at the conflate() boundary

context