You have a fast producer and a slow consumer. Explain how Reactor handles this, and what happens specifically with Flux.interval and Flux.create.
answer
- cold source = waits for demand, no problem
- Flux.interval ignores demand -> overflow IllegalStateException
- fix interval with onBackpressureDrop/Latest
- Flux.create(sink, OverflowStrategy): BUFFER/DROP/LATEST/ERROR/IGNORE
- Flux.generate = demand-driven (backpressure-free by design)
basics
~20 sFor sources that honor backpressure, the producer simply waits for demand — no problem. But timer sources like Flux.interval can't slow down and will overflow (default onError) if the consumer lags; Flux.create needs an explicit OverflowStrategy.
solid answer
~40 sIt depends on whether the source respects demand. Cold generative sources (Flux.range, fromIterable, generate) pause production until the consumer requests more — backpressure just works. But some sources are inherently push-based and cannot slow down: Flux.interval emits on a timer regardless of demand; if the consumer can't keep up, undelivered ticks overflow and the stream fails with an IllegalStateException ('Could not emit tick ... due to lack of requests') — you must add onBackpressureDrop/Latest/Buffer to survive. Flux.create(sink, OverflowStrategy) takes the strategy directly: BUFFER (default-ish, unbounded), DROP, LATEST, ERROR, or IGNORE (violate the contract, push regardless). The senior move is to identify the non-backpressure-aware source and attach an explicit overflow policy, sizing buffers and choosing drop-vs-latest-vs-fail per business need, rather than assuming Reactor magically throttles the producer.
code
java · 16 linesimport reactor.core.publisher.FluxSink;
// PROBLEM: interval ignores demand; slow consumer -> overflow IllegalStateException
Flux.interval(Duration.ofMillis(1))
.concatMap(i -> slowProcess(i)) // slower than 1ms -> "Could not emit tick"
.subscribe();
// FIX: shed late ticks explicitly
Flux.interval(Duration.ofMillis(1))
.onBackpressureDrop(dropped -> log.debug("dropped tick {}", dropped))
.concatMap(i -> slowProcess(i))
.subscribe();
// Push source declares its overflow policy up front
Flux<Event> events = Flux.create(sink -> register(e -> sink.next(e)),
FluxSink.OverflowStrategy.LATEST);go deeper
Know that cold sources wait for demand, and timer sources like interval can overflow if the consumer is slow.
Name the fix (onBackpressureDrop/Latest on interval) and that Flux.create takes an OverflowStrategy.
Categorize sources as demand-aware vs push, list all OverflowStrategy values, and choose a policy per business need; contrast create vs generate.
Design end-to-end policies (buffer sizing, fail-fast vs shed) across thread boundaries and integrate with monitoring/alerting on drops and overflow.
## Two categories of source The entire question hinges on **whether the source can honor demand**. ### Backpressure-aware (cold, generative) sources `Flux.range`, `Flux.fromIterable`, `Flux.generate`, `Flux.fromStream`, R2DBC/reactive-driver streams. These generate items lazily, one demand-batch at a time. If the consumer is slow, the source simply **doesn't generate** the next item until `request(n)` arrives. Fast-producer/slow-consumer is a non-issue: the producer's speed is gated by demand. Memory stays bounded automatically. ### Non-backpressure-aware (hot, timer, or push) sources `Flux.interval`, `Flux.create` (in push modes), broker listeners, external callbacks. These produce on **their own schedule** and physically cannot pause to match demand. Here fast-producer/slow-consumer causes **overflow**, and you must decide the policy. ## Flux.interval specifically `Flux.interval(Duration.ofMillis(1))` emits an incrementing `Long` every interval on the `Schedulers.parallel()` timer, **ignoring downstream demand**. If the consumer processes slower than the tick rate, ticks accumulate with no demand and the stream terminates with an error — Reactor throws an overflow `IllegalStateException` whose message is along the lines of *"Could not emit tick N due to lack of requests (interval doesn't support small downstream requests that replenish slower than the ticks)"*. **Fix**: attach an overflow operator, e.g. `Flux.interval(...).onBackpressureDrop()` or `.onBackpressureLatest()`, so late ticks are shed instead of blowing up. The choice depends on whether you need every tick (you usually don't) or just the latest. ## Flux.create and OverflowStrategy `Flux.create(Consumer<FluxSink<T>> emitter, FluxSink.OverflowStrategy strategy)` lets a push producer emit via `sink.next(...)`. Because the producer isn't demand-driven, you declare how to handle backpressure via **`FluxSink.OverflowStrategy`**: - `BUFFER` — buffer all signals (potentially unbounded → OOM risk). - `DROP` — drop the incoming signal when downstream can't keep up. - `LATEST` — keep only the latest signal. - `ERROR` — signal an `IllegalStateException` on overflow. - `IGNORE` — ignore downstream backpressure entirely; the producer pushes regardless, and downstream operators may then throw overflow errors. `Flux.push(...)` is the single-producer variant with the same overflow semantics. ## Flux.generate vs Flux.create `Flux.generate` is **synchronous and demand-driven** — it produces exactly one item per `request`, so it is backpressure-aware by construction. `Flux.create` is **asynchronous/multi-threaded** and push-oriented, hence the explicit overflow strategy. Choosing `generate` over `create` when you control production is often the cleanest way to get backpressure for free. ## Where thread boundaries fit in Operators like `publishOn(scheduler, prefetch)` insert a bounded queue between producer and consumer threads and request upstream in prefetch-sized batches — they participate in backpressure. But a queue only helps a **backpressure-aware** upstream; it cannot stop a timer from ticking. ## Senior gotchas - Never assume Reactor 'throttles the producer' — it only signals demand; a timer/push source that ignores demand will overflow. - Unbounded `BUFFER` (in `create` or `onBackpressureBuffer`) is the classic OOM. - `onBackpressureDrop` on `Flux.interval` is the idiomatic fix for 'interval overflow' errors. - Test with a deliberately slow consumer (e.g. `.delayElements` downstream, or `StepVerifier` with `thenAwait`) to surface overflow before production does.
- Why does Flux.interval throw when the consumer is slow, but Flux.range does not?Flux.interval emits on a timer regardless of demand, so unconsumed ticks overflow. Flux.range is demand-driven — it only produces the next value when request(n) arrives, so a slow consumer just slows production, never overflows.
- What is the difference between Flux.create and Flux.generate regarding backpressure?Flux.generate is synchronous and demand-driven: exactly one emission per request, so it honors backpressure automatically. Flux.create is asynchronous/push-based and needs an explicit FluxSink.OverflowStrategy because the producer isn't gated by demand.
saying these in an interview costs you the question
- Claiming Reactor automatically slows a Flux.interval producer to match the consumer
- Not knowing Flux.interval overflows with an IllegalStateException
- Thinking Flux.create honors backpressure without specifying an OverflowStrategy
- Confusing Flux.generate (demand-driven) with Flux.create (push)
- Believing publishOn's queue can throttle a timer source