What is backpressure in Project Reactor, and how does request(n) implement it?
answer
- Subscription.request(n) = demand signal
- publisher never emits more than requested
- demand is cumulative + propagates upstream
- request(Long.MAX_VALUE) = no backpressure
- cold sources honor it, timers/create don't
basics
~20 sBackpressure 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 sBackpressure 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// 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
Know the definition: slow consumer controls fast producer via request(n), and the publisher can't exceed requested demand.
Explain cumulative demand, that plain subscribe requests Long.MAX_VALUE, and which sources honor backpressure natively vs need overflow strategies.
Discuss demand propagation/reshaping through operators (buffer, flatMap, limitRate) and the pull-push hybrid model.
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