skip to content

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

level: middleimportance: should knowfreq 55%

answer

  1. caps upstream request size, not time
  2. reshapes request(MAX) into bounded batches
  3. default replenish at 75% consumed
  4. limitRate(highTide, lowTide) to tune
  5. same logic as operator 'prefetch' (default 256)

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.

solid answer

~50 s

limitRate(prefetchRate) sits in a pipeline and splits a large (often unbounded) downstream request into bounded upstream requests of at most prefetchRate items. It requests prefetchRate upfront, then replenishes when a low-water mark is reached — the default is 75% consumed, i.e. it requests the next 75% batch once 75% of the current batch has been delivered (limitRate(highTide, lowTide) lets you tune this; lowTide=0 means wait until fully drained before re-requesting). This protects an upstream that would otherwise be told 'send everything' (request(Long.MAX_VALUE)) — common when the terminal subscriber is unbounded. Typical uses: capping how many rows an R2DBC/database driver fetches at once, throttling a paged API source, or bounding memory of an intermediate stage. It shapes demand, not time — it is not a rate-limiter in requests-per-second; for that use delayElements or a rate-limiting library.

code

java · 10 lines
java
// Unbounded subscriber would pull everything; limitRate bounds upstream to 50 at a time
repository.findAllStreaming()          // R2DBC Flux<Row>, could be huge
    .limitRate(50)                     // request 50, replenish at ~75% (after 37-38 consumed)
    .map(this::toDto)
    .subscribe(this::send);            // even this unbounded subscriber now pulls in bounded chunks

// Stricter: request 50, then wait until ALL 50 drained before the next batch
repository.findAllStreaming()
    .limitRate(50, 0)                  // highTide=50, lowTide=0
    .subscribe(this::send);

go deeper

for a junior

Know it bounds how much is requested upstream at once, protecting memory when the consumer is unbounded.

for a middle

Explain the 75% replenish default, the highTide/lowTide two-arg form, and that it is demand-based not time-based.

for a senior

Relate it to operator prefetch (default 256), and pick sensible batch sizes for DB/remote sources.

for a principal

Reason about round-trip overhead vs memory tradeoffs and where demand-shaping belongs in a multi-stage reactive pipeline.

## The problem Many terminal subscribers request `Long.MAX_VALUE` (unbounded). A cold source will happily try to satisfy that, but you may not want an operator or driver pulling an unbounded amount into memory at once. **`limitRate`** re-shapes that unbounded downstream demand into bounded upstream requests. ## What limitRate(n) does `Flux<T> limitRate(int prefetchRate)`: - It requests **at most `prefetchRate`** items from upstream at a time, regardless of how much its downstream requested. - It maintains an internal bounded queue of size `prefetchRate` and relays items downstream as demand allows. - As the batch drains, it **replenishes**: by default it re-requests when **75%** of the batch has been consumed (a low-water/high-water design). So with `limitRate(100)` it requests 100, and after 75 have been produced downstream it requests the next 75 to top back up. This avoids stalls (it refills before running dry) while keeping the outstanding amount bounded. ## Tuning with the two-arg form `limitRate(int highTide, int lowTide)`: - `highTide` = the batch/prefetch amount. - `lowTide` = the replenish threshold expressed as the number already consumed. The default single-arg form is equivalent to `lowTide = highTide * 0.75` (75% replenish). - `lowTide = 0` means **wait until the whole batch is drained** before requesting the next — stricter, fewer overlapping requests. - `lowTide = highTide` would replenish immediately after each item (effectively continuous topping-up). ## Relationship to prefetch Many operators (`flatMap`, `concatMap`, `publishOn`) already take a **prefetch** parameter that does the same bounded-request-with-75%-replenish internally (default prefetch is 256, `Queues.SMALL_BUFFER_SIZE`). `limitRate` is the standalone operator form of that same demand-shaping logic, useful when you want to bound a stage that doesn't expose its own prefetch, or override the default. ## When to use it - **Bounding a database/driver fetch**: an R2DBC or reactive Mongo stream feeding an unbounded subscriber — `limitRate` caps how many rows are pulled per round trip, controlling memory. - **Paged/remote sources**: throttle how aggressively you pull pages. - **Memory control** on an intermediate transform that buffers. ## Critical clarification — it is NOT a time-based rate limiter `limitRate` bounds **demand quantity**, not throughput per unit time. It will not slow a source to 'N per second'. If you want temporal throttling, use `delayElements(Duration)`, a windowing approach, or an external rate limiter (e.g. Resilience4j / Bucket4j). Confusing the two is a common interview trap. ## Gotchas - If the source honors backpressure and the consumer is already bounded, `limitRate` may add little — it matters most when the downstream demand is unbounded. - Setting `prefetchRate` too small increases the number of upstream request round-trips (overhead); too large increases memory. Tune to the workload.

  • Is limitRate a requests-per-second throttle?
    No. It bounds the *quantity* of demand (how many items are requested upstream at once), not throughput over time. For per-second throttling use delayElements or an external rate limiter.
  • What does the second argument in limitRate(highTide, lowTide) control?
    The replenish threshold: after lowTide items of the current highTide batch have been consumed, it requests the next batch. lowTide=0 waits for full drain; the single-arg form defaults lowTide to 75% of highTide.

saying these in an interview costs you the question

  • Describing limitRate as a per-second/time-based rate limiter
  • Thinking it changes what the downstream subscriber receives (it changes upstream request pattern, not delivered content)
  • Not knowing the default replenishment happens at 75%
  • Confusing it with delayElements

context