skip to content

How is backpressure handled on a result Flux from a reactive MongoDB query?

level: seniorimportance: must knowfreq 45%

answer

  1. request(n) -> driver getMore per batch
  2. Demand reaches the DB cursor
  3. cursorBatchSize tunes granularity
  4. limitRate / bounded flatMap bound in-flight
  5. Hot change stream needs onBackpressure*; collectList kills streaming

basics

~20 s

The reactive driver honours reactive-streams demand: the subscriber requests N items, and the driver fetches from MongoDB in batches to match, pausing when demand is zero. So a slow consumer naturally throttles the query instead of buffering everything.

solid answer

~50 s

A result Flux is a reactive-streams Publisher, so it respects downstream demand. The subscriber signals request(n); the MongoDB reactive driver translates that into cursor batch fetches (getMore) and only pulls more documents as demand arrives, pausing when demand is zero. This means backpressure propagates all the way to the database cursor, so a slow consumer (for example a slow HTTP client draining an SSE stream) throttles fetching rather than loading the whole collection into memory. You tune the flow with the cursor batch size and with operators: limitRate(n) caps outstanding demand, onBackpressureBuffer/Drop/Latest set overflow strategy when you insert an async boundary the source cannot slow (like a hot change stream), and buffer/window batch elements. The natural request-driven model is the default; you only add explicit strategies when a fast, non-throttleable producer meets a slow consumer.

code

java · 12 lines
java
// Stream a large collection to a slow SSE client with bounded in-flight work.
@GetMapping(value="/export", produces=MediaType.TEXT_EVENT_STREAM_VALUE)
Flux<Dto> export() {
    return repo.findAll()               // cold cursor: getMore only on demand
        .limitRate(100)                 // cap outstanding request to 100
        .flatMap(this::enrich, 8)       // bounded concurrency, not unbounded
        .onBackpressureBuffer(1_000);   // guard an async boundary
}

// Query-level batch size (finer vs coarser backpressure granularity):
Query q = new Query().cursorBatchSize(200);
Flux<User> users = template.find(q, User.class);

go deeper

for a junior

Know backpressure means the consumer controls the pace so it is not overwhelmed.

for a middle

Explain request(n) driving cursor getMore and that demand reaches the DB.

for a senior

Tune with cursorBatchSize, limitRate, bounded flatMap, and pick overflow strategies for hot sources.

for a principal

Reason about cold vs hot producers, async-boundary buffering, memory bounds, and end-to-end flow control under load.

**Backpressure** is the reactive-streams mechanism by which a **subscriber controls how fast a publisher produces**, so a fast source cannot overwhelm a slow consumer's memory. It is expressed through the `Subscription.request(n)` signal: the subscriber asks for at most `n` more elements; the publisher must not emit more than the outstanding demand. **How a Mongo result Flux honours it.** The **MongoDB Reactive Streams driver** is genuinely demand-aware. When you subscribe to a `Flux` from `findAll()` or `template.find(...)`: 1. The driver opens a **cursor** and fetches the **first batch**. 2. As the subscriber calls `request(n)`, the driver emits buffered documents and, when the batch is drained and more demand exists, issues a **`getMore`** to fetch the next batch. 3. If downstream demand is **zero**, the driver stops issuing `getMore` and the cursor simply waits. Nothing is fetched ahead of demand beyond one batch. So backpressure propagates **all the way to the database cursor**: a slow consumer throttles the query itself, and you never materialise the entire result set in memory. This is the core reason reactive Mongo streams huge result sets safely. **Batch size.** The **cursor batch size** (via `Query.cursorBatchSize(n)` / template options, mapping to the driver's `batchSize`) controls the granularity: how many documents each network round trip fetches. Smaller batches = finer backpressure granularity, more round trips; larger batches = fewer round trips but coarser flow control and more memory per batch. **Operators for shaping demand.** - **`limitRate(n)`** — caps the outstanding request to `n`, prefetching in chunks; useful to bound in-flight documents even if a downstream operator requests `Long.MAX_VALUE` (unbounded). - **`buffer(n)` / `window(n)` / `bufferTimeout`** — group elements to process in batches (e.g., bulk write downstream). - **`onBackpressureBuffer` / `onBackpressureDrop` / `onBackpressureLatest`** — overflow strategies for when the **source cannot be slowed** (a *hot* publisher). A plain query Flux is *cold* and throttleable, so you rarely need these; but a **change stream** or **`@Tailable`** feed of live inserts is effectively push-based — if events arrive faster than the consumer drains, you choose to buffer (bounded, or risk OOM), drop oldest/newest, or keep only the latest. **Where backpressure breaks down (gotchas).** - **Unbounded request upstream.** Some operators/subscribers request `Long.MAX_VALUE`, signalling "give me everything." The cold cursor still only fetches per batch, but any intermediate **async boundary** (`publishOn`, `flatMap` with high concurrency) can buffer unboundedly. Use `limitRate` / bounded `flatMap` concurrency. - **`flatMap` concurrency.** `flatMap` has a default concurrency (256) and prefetch; a slow inner publisher plus a fast outer source can accumulate in-flight work. Bound it: `flatMap(fn, concurrency)`. - **Hot sources.** Change streams / tailable feeds do not slow the database (other writers keep inserting); if your consumer lags you *must* pick an overflow policy or you buffer without bound and risk OOM. - **`.collectList()` defeats streaming.** Collecting a huge Flux into a `List<T>` (`Mono<List<T>>`) buffers the entire result in memory, discarding the backpressure benefit. Stream to the client instead. - **Blocking the consumer.** A blocking operation on the pipeline stalls demand and can starve the event loop; keep the consumer non-blocking. **When to add explicit strategies.** Default request-driven flow is enough for ordinary queries. Add `limitRate`, bounded `flatMap`, `buffer`, or `onBackpressure*` when you (a) fan out to slower downstream work, (b) bridge to a hot producer (change stream), or (c) need to bound memory under an unbounded-request subscriber.

  • Does calling limitRate change how MongoDB fetches, or only the operator chain?
    limitRate bounds the demand the operator propagates upstream. Because the reactive driver ties getMore to demand, capping demand does reduce how aggressively the cursor prefetches, so it can influence fetching, not just the in-memory operator queue.
  • Why can a change-stream Flux overflow when a plain findAll Flux does not?
    findAll is a cold cursor: with no demand it simply stops issuing getMore, so it self-throttles. A change stream is effectively hot: other clients keep writing and events keep arriving regardless of your demand, so a lagging consumer must buffer or drop, otherwise memory grows unbounded.
  • What is the downside of collectList() on a large query?
    It buffers the entire result set into a single in-memory List before emitting, so peak memory scales with result size and you lose streaming/backpressure. Prefer streaming the Flux element-by-element.

saying these in an interview costs you the question

  • Believing the driver always fetches the entire result set into memory regardless of demand
  • Thinking backpressure stops at the operator chain and never reaches the DB cursor
  • Assuming a change-stream/tailable hot source self-throttles like a normal query
  • Using collectList and still claiming the pipeline is backpressure-friendly

context