skip to content

Reactive Streams & Project Reactor

The Reactive Streams contract and the Reactor library that implements it: Mono and Flux, the operator vocabulary, schedulers, backpressure and error handling. Nothing in WebFlux makes sense until this layer does, so interviews start here.

part ofSpring Frameworkoverview, primer and where to startread it →
on this pageshow

explore

questions

30

What is backpressure in Project Reactor, and how does request(n) implement it?

level: juniorimportance: must knowfreq 78%

answer

  1. Subscription.request(n) = demand signal
  2. publisher never emits more than requested
  3. demand is cumulative + propagates upstream
  4. request(Long.MAX_VALUE) = no backpressure
  5. cold sources honor it, timers/create don't

basics

~20 s

Backpressure lets a slow consumer control how fast a producer sends data. The subscriber calls request(n) on its Subscription to ask for at most n items; the publisher must never emit more than was requested.

solid answer

~40 s

Backpressure is the Reactive Streams flow-control mechanism that stops a fast producer from overwhelming a slow consumer. It is pull-driven: when a Subscriber subscribes, the Publisher hands it a Subscription, and the Subscriber calls subscription.request(n) to signal demand for n elements. The Publisher must never emit more than the total requested. Demand accumulates, so a consumer can request(1) at a time or request more as it catches up. Operators propagate demand upstream — each requests from its own upstream based on what its downstream asked for. Most cold sources (Flux.range, Flux.fromIterable, R2DBC, WebClient bodies) honor this natively. When you block() or subscribe without limiting, demand is effectively request(Long.MAX_VALUE) — 'give me everything' — which disables backpressure.

code

java · 15 lines
java
// A subscriber that pulls one item at a time (explicit backpressure)
Flux.range(1, 100)
    .doOnRequest(n -> System.out.println("upstream asked for " + n))
    .subscribe(new BaseSubscriber<Integer>() {
        @Override
        protected void hookOnSubscribe(Subscription subscription) {
            request(1); // signal demand for exactly one element
        }
        @Override
        protected void hookOnNext(Integer value) {
            System.out.println("got " + value);
            request(1); // ask for the next only after processing this one
        }
    });
// Contrast: subscribe(System.out::println) would request(Long.MAX_VALUE) — no backpressure.

go deeper

for a junior

Know the definition: slow consumer controls fast producer via request(n), and the publisher can't exceed requested demand.

for a middle

Explain cumulative demand, that plain subscribe requests Long.MAX_VALUE, and which sources honor backpressure natively vs need overflow strategies.

for a senior

Discuss demand propagation/reshaping through operators (buffer, flatMap, limitRate) and the pull-push hybrid model.

for a principal

Tie demand signaling to end-to-end flow control across schedulers and over Reactor Netty/TCP, and reason about where unbounded demand leaks in.

## The problem backpressure solves In a streaming pipeline a **producer** (`Publisher`) generates data and a **consumer** (`Subscriber`) processes it. If the producer is faster than the consumer, items pile up in memory and you eventually hit an `OutOfMemoryError`. **Backpressure** is the mechanism that lets the consumer tell the producer 'only send what I can handle'. ## The Reactive Streams contract Project Reactor implements the **Reactive Streams** specification. Its interfaces are `Publisher`, `Subscriber`, `Subscription`, and `Processor`. The flow is: 1. You call `flux.subscribe(subscriber)`. 2. The `Publisher` calls `subscriber.onSubscribe(subscription)`, handing the subscriber a **`Subscription`** handle. 3. Nothing is emitted yet. The subscriber must call `subscription.request(n)` to signal **demand** — it wants up to `n` elements. 4. The publisher then invokes `onNext(item)` at most `n` times, then waits for more demand. 5. When done it calls `onComplete()`; on failure `onError(throwable)`. **Demand is cumulative**: `request(2)` then `request(3)` means the publisher may emit up to 5. The publisher is contractually forbidden from ever emitting more than the outstanding demand — that is the core invariant of backpressure. ## Pull-push hybrid Reactive Streams is often called a 'dynamic pull-push' model. When demand is high the publisher pushes freely; when demand is low the consumer effectively pulls one at a time. `request(Long.MAX_VALUE)` means 'unbounded — send everything as fast as you can', which turns off backpressure. Blocking operators like `block()`, `toIterable()`, and simple `subscribe()` (with no custom subscriber) request unbounded. ## Demand propagation through operators Each operator is both a subscriber (to its upstream) and a publisher (to its downstream). It relays demand: if downstream requests 10, a pass-through operator like `map` requests 10 from upstream. Some operators reshape demand — `buffer(5)` requests 5 upstream items per 1 downstream request; `flatMap` requests based on its concurrency; `limitRate(n)` caps an unbounded downstream request into bounded upstream chunks. ## Who honors backpressure? - **Cold, generative sources** — `Flux.range`, `Flux.fromIterable`, `Flux.generate`, R2DBC result streams, `WebClient` response bodies — honor demand natively because they can pause generation. - **Time/event-driven sources** — `Flux.interval`, `Flux.create` (push), UI events, message brokers — produce on their own schedule and **cannot** naturally slow down; they need an explicit overflow strategy (`onBackpressureBuffer/Drop/Latest/Error`). ## Gotchas - `subscribe()` with no arguments requests `Long.MAX_VALUE`: no backpressure. To exert it, use `BaseSubscriber` and call `request()` yourself, or an operator like `limitRate`. - Backpressure is about **demand signaling**, not thread throttling — it works across `publishOn`/`subscribeOn` thread boundaries because demand flows through the `Subscription`. - Over the network (Reactor Netty) demand is ultimately backed by TCP flow control.

  • What demand does a plain flux.subscribe(System.out::println) issue?
    Long.MAX_VALUE (unbounded). The default LambdaSubscriber requests everything, so backpressure is effectively disabled and the source emits as fast as it can.
  • Does backpressure still work when you cross threads with publishOn?
    Yes. Demand flows through the Subscription regardless of scheduler. publishOn actually introduces a bounded queue (default prefetch 256) and requests upstream in batches, so it participates in backpressure rather than breaking it.

saying these in an interview costs you the question

  • Thinking backpressure means the framework automatically slows the producer with no demand signaling
  • Believing a plain subscribe() applies backpressure
  • Confusing backpressure with thread pool size or rate limiting
  • Claiming all sources (including Flux.interval) honor backpressure natively

context

open as a page

In Project Reactor, what happens to a Flux or Mono when an error occurs, and why can't you just use a try/catch around it?

level: juniorimportance: must knowfreq 70%

basics

~20 s

An error is a terminal signal: the stream emits onError, stops, and sends nothing more. try/catch doesn't work because the data flows asynchronously later, not on the line that builds the pipeline. You recover with operators like onErrorResume.

open as a page

What is the difference between Mono and Flux in Project Reactor, and when would you use each?

level: juniorimportance: must knowfreq 85%

basics

~20 s

Both are reactive publishers. Mono emits 0 or 1 item then completes; Flux emits 0 to many items then completes. Use Mono for a single result (like one entity), Flux for a stream or collection of items.

open as a page

In Project Reactor, what is the difference between map and flatMap on a Flux or Mono?

level: juniorimportance: must knowfreq 85%

basics

~10 s

map transforms each item synchronously, one value in -> one value out. flatMap transforms each item into a Publisher (another Flux/Mono, often async) and flattens all those inner streams into one output stream.

open as a page

What are the four interfaces defined by the Reactive Streams specification, and what is each responsible for?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Publisher produces data, Subscriber consumes it, Subscription links the two and lets the Subscriber request items or cancel, and Processor is both a Subscriber and a Publisher (a middle stage).

open as a page

What is a Reactor Scheduler, and what are the four built-in Schedulers (boundedElastic, parallel, single, immediate) each meant for?

level: juniorimportance: must knowfreq 70%

basics

~10 s

A Scheduler decides which thread(s) run reactive work. boundedElastic is for blocking I/O, parallel for CPU work, single for one shared thread, and immediate runs on the current thread (no switch).

open as a page

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

level: middleimportance: must knowfreq 72%

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.

open as a page

Compare onErrorReturn, onErrorResume, and onErrorMap. When would you choose each?

level: middleimportance: must knowfreq 75%

basics

~10 s

onErrorReturn emits a fixed fallback value then completes. onErrorResume switches to another Mono/Flux (e.g., a cache or default stream). onErrorMap translates the exception into a different one but the stream still ends in error.

open as a page

Walk through the main Reactor factory methods (just, empty, error, fromCallable, defer, fromIterable) and when you'd reach for each.

level: middleimportance: must knowfreq 75%

basics

~20 s

just wraps a ready value; empty completes with nothing; error emits a failure; fromCallable defers a synchronous/blocking call and captures its exceptions; defer builds a fresh publisher lazily per subscription; fromIterable turns a collection into a Flux.

open as a page

Contrast flatMap, concatMap, and flatMapSequential in terms of ordering and concurrency. When would you pick each?

level: middleimportance: must knowfreq 78%

basics

~10 s

flatMap: concurrent inner subscriptions, output order not preserved. concatMap: one inner at a time (sequential), order preserved, no concurrency. flatMapSequential: concurrent inner subscriptions but output re-ordered to match source order.

open as a page

What do Subscription.request(n) and Subscription.cancel() do, and what are the rules around calling them?

level: middleimportance: must knowfreq 60%

basics

~20 s

request(n) tells the Publisher the Subscriber is ready to receive up to n more elements — it is the demand signal. cancel() tells the Publisher to stop emitting and release resources. Both are called by the Subscriber on its Subscription.

open as a page

Describe the signal protocol a Subscriber receives: onSubscribe, onNext, onError, onComplete — their ordering and cardinality.

level: middleimportance: must knowfreq 65%

basics

~10 s

onSubscribe is called first, exactly once. Then onNext happens zero or more times. The stream ends with exactly one terminal signal: either onError or onComplete — never both, and nothing comes after it.

open as a page

What is the difference between subscribeOn and publishOn, and how does each affect which thread the operators run on?

level: middleimportance: must knowfreq 85%

basics

~10 s

publishOn switches the thread for operators placed after it (downstream). subscribeOn sets the thread for the source and the whole subscription, no matter where you put it in the chain.

open as a page

You must call a legacy blocking library (JDBC / RestTemplate) inside a WebFlux handler. How do you keep from blocking the event loop, and what does the correct code look like?

level: seniorimportance: must knowfreq 80%

basics

~10 s

Wrap the blocking call in Mono.fromCallable(...) (or Flux) and add subscribeOn(Schedulers.boundedElastic()), so the blocking work runs on the I/O pool instead of the Netty event-loop thread.

open as a page

Explain the difference between assembly time and subscription time in Reactor. Why does it matter for correctness, and how do defer/fromCallable relate to it?

level: principalimportance: must knowfreq 55%

basics

~20 s

Assembly time is when you build the operator chain (the reactive pipeline is just a blueprint). Subscription time is when someone calls subscribe() and data actually flows. Nothing runs until subscription. Mono.defer/Mono.fromCallable delay eager work so it happens per-subscription instead of at assembly.

open as a page

What does limitRate(n) do, and how does it reshape demand from an unbounded downstream?

level: middleimportance: should knowfreq 55%

basics

~20 s

limitRate(n) caps how much an operator requests from its upstream at once. Even if the downstream asks for everything, limitRate requests only n at a time, replenishing as items are consumed — so upstream sees bounded demand.

open as a page

How do you convert between Mono and Flux? Explain flatMapMany, next, and collectList and when each applies.

level: middleimportance: should knowfreq 60%

basics

~20 s

Mono to Flux: use flatMapMany when the single item expands into many (returns a Flux), or Flux.from(mono). Flux to Mono: next() takes the first element (0..1), collectList() gathers all elements into a Mono<List>, and single()/last() take exactly one/the last.

open as a page

You have a fast producer and a slow consumer. Explain how Reactor handles this, and what happens specifically with Flux.interval and Flux.create.

level: seniorimportance: should knowfreq 58%

basics

~20 s

For 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.

open as a page

What does onErrorContinue do, why is it considered controversial, and how does it differ from onErrorResume inside a flatMap?

level: seniorimportance: should knowfreq 45%

basics

~20 s

onErrorContinue drops the element that caused an error and keeps the stream going instead of terminating. It's controversial because it breaks the 'error is terminal' model and only works with operators that explicitly support it, so behavior is surprising. Prefer onErrorResume inside flatMap.

open as a page

How does retryWhen with Retry.backoff work, and how does it differ from a plain retry()? What happens when attempts are exhausted?

level: seniorimportance: should knowfreq 65%

basics

~10 s

retry() re-subscribes to the source immediately on error, up to N times. retryWhen(Retry.backoff(max, minDelay)) re-subscribes with exponential backoff plus jitter. When retries run out it fails with a RetryExhaustedException wrapping the last error.

open as a page

Explain why Mono.just(Instant.now()) gives every subscriber the same value while Mono.defer(() -> Mono.just(Instant.now())) does not. What does this reveal about assembly time vs subscription time?

level: seniorimportance: should knowfreq 55%

basics

~20 s

just captures its value once when the pipeline is built (assembly time), so all subscribers share it. defer runs its supplier on each subscription, producing a fresh value per subscriber. It shows Reactor separates building the pipeline from executing it.

open as a page

What do buffer and window do, and how does transform differ from transformDeferred? Give a use case for each.

level: seniorimportance: should knowfreq 40%

basics

~20 s

buffer groups incoming elements into Lists (by size or time) and emits Flux<List<T>>. window is similar but emits Flux<Flux<T>> — each group is itself a stream. transform lets you factor a reusable chain of operators into a function applied once at assembly time; transformDeferred applies it per subscriber.

open as a page

Explain zip, merge, and switchMap. How do they combine or select among streams, and what are their typical use cases?

level: seniorimportance: should knowfreq 62%

basics

~20 s

zip pairs one item from each source into a combined tuple, emitting when all have a value. merge interleaves items from multiple sources concurrently as they arrive. switchMap maps each element to a Publisher but cancels the previous inner when a new element arrives, keeping only the latest.

open as a page

Why does the Reactive Streams specification exist, how does Project Reactor relate to it, and what is its relationship to java.util.concurrent.Flow?

level: seniorimportance: should knowfreq 45%

basics

~10 s

The spec gives asynchronous stream libraries a common, vendor-neutral contract so they interoperate with built-in flow control. Reactor implements it (Flux/Mono are Publishers). Java 9's java.util.concurrent.Flow contains the identical interfaces, mirrored into the JDK.

open as a page

Why is Schedulers.parallel() the wrong place for blocking calls, and boundedElastic the wrong place for tight CPU-bound work? What can go wrong?

level: seniorimportance: should knowfreq 55%

basics

~20 s

parallel() has only one thread per CPU core, so blocking those threads starves the pool and can deadlock. boundedElastic can hold many threads, so putting CPU work there causes oversubscription and context-switching instead of speedup.

open as a page

In a Spring WebFlux app on Reactor Netty, how does backpressure propagate end-to-end from the HTTP client, across the network, to your reactive pipeline?

level: principalimportance: should knowfreq 40%

basics

~20 s

Reactor Netty maps reactive demand onto TCP flow control. When your response Publisher is slow (or a slow client stops reading), Netty stops requesting/writing, its socket buffers fill, the TCP receive window shrinks, and the sender is throttled — backpressure crosses the network via TCP.

open as a page

Design the error-handling layer for a WebFlux endpoint that aggregates two downstream services. Discuss ordering of doOnError/retry/onErrorMap/onErrorResume, and how error signals interact with context propagation.

level: principalimportance: should knowfreq 35%

basics

~20 s

Scope handling per call: retry transient errors with backoff around each downstream, log with doOnError, translate infra exceptions with onErrorMap at the boundary, then provide a fallback with onErrorResume. Keep operator order deliberate because each catches only upstream errors, and rely on the Reactor Context to carry request/trace data through the error path.

open as a page

As a tech lead reviewing a WebFlux codebase, how do Mono/Flux cardinality and the eager-vs-lazy factory choices shape API contracts, error handling, and blocking-code integration? What patterns do you enforce?

level: principalimportance: should knowfreq 35%

basics

~20 s

Type return values by true cardinality (Mono for 0..1, Flux for 0..N) so contracts are honest. Ban Mono.just around blocking/side-effecting calls; require fromCallable+subscribeOn(boundedElastic) or defer for retriable/fresh work. Push errors into onError, avoid collectList on unbounded streams, and never block the event loop.

open as a page

What contractual guarantees must a compliant Publisher provide to its Subscriber, and why do these rules matter for building operators?

level: principalimportance: should knowfreq 30%

basics

~20 s

A compliant Publisher must: call onSubscribe first and exactly once, never emit more onNext than requested, deliver signals serially (never concurrently), end with at most one terminal signal, never emit after a terminal or after cancel, and never pass null. These guarantees let operators compose safely.

open as a page

Explain why subscribeOn is position-independent while publishOn is not, and how you'd reason about which thread each operator runs on in a chain that mixes both.

level: principalimportance: should knowfreq 40%

basics

~20 s

subscribeOn hooks the subscription, which flows upward to the source, so its position doesn't matter — it sets where the source runs. publishOn hooks the downward data signals, switching threads only for operators below it, so position matters.

open as a page