skip to content

Compare onBackpressureBuffer, onBackpressureDrop, onBackpressureLatest, and onBackpressureError. When do you reach for each?

level: middleimportance: must knowfreq 72%

answer

  1. Buffer = queue (unbounded = OOM risk)
  2. Drop = discard newcomers
  3. Latest = keep only freshest
  4. Error = fail fast with IllegalStateException
  5. BufferOverflowStrategy: ERROR/DROP_LATEST/DROP_OLDEST

basics

~10 s

They decide what happens when a source emits faster than the consumer requests: Buffer queues the extras, Drop discards new ones, Latest keeps only the newest, and Error fails with an overflow exception.

solid answer

~40 s

These operators bridge a non-backpressure-aware source to a slow consumer by choosing an overflow policy. onBackpressureBuffer queues unrequested items (unbounded by default, or bounded with a max size + BufferOverflowStrategy callback) — use when you can't lose data and bursts are temporary. onBackpressureDrop silently discards items that arrive with no outstanding demand (optional onDropped callback) — use for lossy telemetry/sensor streams where newest-ish is fine. onBackpressureLatest keeps only the most recent item and drops older undelivered ones — use for 'current state' feeds like live prices or a UI gauge. onBackpressureError signals an IllegalStateException (overflow) the moment the buffer would be exceeded — use fail-fast when overflow indicates a real bug. The right choice is a data-loss vs memory-risk vs fail-fast tradeoff; unbounded Buffer is the hidden OOM footgun.

code

java · 17 lines
java
import reactor.core.publisher.BufferOverflowStrategy;

// Fast push source adapted to a slow consumer, four ways:
Flux<Tick> ticks = fastTickSource();

// 1) Bounded buffer, drop oldest when full, log overflow
ticks.onBackpressureBuffer(1000, dropped -> log.warn("dropped {}", dropped),
                           BufferOverflowStrategy.DROP_OLDEST);

// 2) Drop new items when no demand, count them
ticks.onBackpressureDrop(t -> meter.increment("ticks.dropped"));

// 3) Keep only the freshest value (live gauge)
ticks.onBackpressureLatest();

// 4) Fail fast — emits onError(IllegalStateException) on overflow
ticks.onBackpressureError();

go deeper

for a junior

Name the four and their one-line behavior (queue / drop new / keep latest / fail).

for a middle

Explain the data-loss vs memory vs fail-fast tradeoff, the unbounded-buffer OOM risk, and pick the right one per scenario.

for a senior

Discuss BufferOverflowStrategy (ERROR/DROP_LATEST/DROP_OLDEST), operator placement as a producer/consumer boundary, and observability callbacks.

for a principal

Reason about combining these with prefetch/publishOn queues and when to fail-fast vs degrade, plus SLA implications of silent drops.

## Why these operators exist Some sources cannot slow down on demand — a `Flux.create` push sink, a message-broker listener, mouse events, a metrics stream. When downstream demand is exhausted but the source keeps producing, you have **overflow**. The `onBackpressure*` operators sit in the pipeline and define the overflow policy. They effectively act as an adapter between an unbounded upstream and a backpressure-respecting downstream. ## The four strategies ### `onBackpressureBuffer()` Queues items that exceed current demand and replays them as demand arrives. - **Default is UNBOUNDED** — the queue grows without limit. If the consumer never catches up, this is an `OutOfMemoryError` waiting to happen. - Overloads: `onBackpressureBuffer(int maxSize)`, `onBackpressureBuffer(int maxSize, Consumer<T> onOverflow)`, and `onBackpressureBuffer(int maxSize, BufferOverflowStrategy strategy)`. - **`BufferOverflowStrategy`** enum: `ERROR` (fail when full), `DROP_LATEST` (drop the incoming item), `DROP_OLDEST` (evict the oldest queued item to make room). - **Use when**: no data may be lost and overload is transient (bursty traffic that drains). ### `onBackpressureDrop()` Discards any element that arrives while downstream demand is zero. Overload `onBackpressureDrop(Consumer<T> onDropped)` lets you count/log drops. - **Use when**: loss is acceptable and you'd rather stay real-time — e.g. high-frequency sensor readings, non-critical telemetry. ### `onBackpressureLatest()` Keeps a single-slot 'latest value'. If new items arrive before the old undelivered one is consumed, the old one is overwritten. The consumer always gets the freshest value when it next requests. - **Use when**: only current state matters — live stock price, temperature gauge, a progress percentage. ### `onBackpressureError()` The instant an item can't be delivered (no demand), it terminates the stream with `onError` carrying an `IllegalStateException` (Reactor's overflow exception, from `Exceptions.failWithOverflow()`). - **Use when**: overflow should never happen and you want to fail fast rather than silently drop or grow memory. ## Placement matters These operators establish a boundary: **upstream of the operator** items are produced freely; **downstream** they are delivered per demand. So `source.onBackpressureLatest().publishOn(scheduler)` means the latest-wins buffer absorbs the mismatch between the fast source and the slower `publishOn` consumer. Putting the operator in the wrong place (e.g. after the slow stage) won't protect the fast stage. ## Common gotchas - **Unbounded `onBackpressureBuffer()` is the #1 footgun** — always bound it in production or ensure the consumer provably keeps up. - These operators are **no-ops for sources that already honor backpressure** (like `Flux.range`) because such sources never overflow — they wait for demand instead. - `onBackpressureError`'s exception is an `IllegalStateException`; don't confuse it with a generic error — check for the overflow marker. - Drop/Latest **silently lose data by design** — always attach the callback for observability, or you'll debug 'missing events' blind.

  • What is the default capacity of onBackpressureBuffer() with no arguments, and why is that dangerous?
    It is unbounded. If the consumer never catches up, the internal queue grows without limit and eventually causes an OutOfMemoryError. Always bound it (maxSize + BufferOverflowStrategy) in production.
  • You have a live price ticker feeding a dashboard that only cares about the current price. Which operator, and why?
    onBackpressureLatest — stale intermediate prices are worthless, so overwrite them and always deliver the freshest value when the consumer requests. Buffering would deliver outdated prices; Drop might drop the newest.

saying these in an interview costs you the question

  • Claiming onBackpressureBuffer() is bounded by default
  • Thinking onBackpressureLatest keeps the oldest item
  • Saying onBackpressureDrop drops the currently-buffered oldest item (it drops incoming ones)
  • Believing these operators magically slow a source that already honors backpressure (they're no-ops there)
  • Not knowing onBackpressureError raises IllegalStateException

context