skip to content

How does backpressure propagate in Reactor Kafka, and how is it implemented in terms of consumer pause/resume?

level: seniorimportance: should knowfreq 40%

answer

  1. request(n) demand -> pause()/resume() on partitions
  2. Kafka fetch is push-batch, no native pull
  3. pause stops fetch but keeps heartbeats
  4. ignore it -> OOM or max.poll.interval.ms eviction
  5. tune max.poll.records / limitRate

basics

~20 s

When a downstream subscriber requests fewer records than are available, Reactor Kafka stops the KafkaConsumer from fetching more by calling consumer.pause() on its partitions. When demand returns, it calls consumer.resume(). This keeps fast brokers from overwhelming slow processing.

solid answer

~50 s

Reactive Streams backpressure means a subscriber signals demand (request(n)) and a publisher must not emit more than requested. Reactor Kafka maps this onto the KafkaConsumer: the receive Flux buffers fetched records, and when downstream demand is exhausted (the buffer fills past a threshold), it invokes consumer.pause() on the assigned partitions so poll() returns no new records but still sends heartbeats. When the downstream catches up and requests more, it calls consumer.resume() and polling continues. This is essential because Kafka's fetch protocol has no per-record pull — the broker pushes batches as fast as the network allows — so pause/resume is the only lever to throttle ingestion. Tuning levers include max.poll.records (cap per poll), and Reactor's prefetch/request sizing. The danger if you ignore backpressure: unbounded in-memory buffering and OOM, or exceeding max.poll.interval.ms (because you stopped polling to process) and getting kicked from the group.

go deeper

for a junior

Know backpressure means a slow consumer can tell the source to slow down so it isn't flooded.

for a middle

Explain that Reactor Kafka uses consumer.pause()/resume() to throttle fetching based on downstream demand.

for a senior

Detail the queue/threshold mechanism, why heartbeats survive pause, and max.poll.interval.ms / OOM failure modes.

for a principal

Design end-to-end flow control across reactive stages, set max.poll.records/limitRate policy, and reason about rebalance safety.

## What backpressure is In **Reactive Streams**, data flows from a Publisher to a Subscriber, but the *Subscriber controls the rate* via `Subscription.request(n)` — it asks for at most n items. A compliant Publisher must never emit more than the outstanding demand. This prevents a fast producer from flooding a slow consumer's memory. ## The Kafka mismatch Kafka's consumer **fetch protocol is push-batch oriented**: `consumer.poll()` returns whatever batches the broker has ready, bounded only by `max.poll.records` and fetch byte limits. There is no native 'give me exactly 3 records' pull. So Reactor Kafka must *bridge* demand-based backpressure onto a batch-fetch client. ## How Reactor Kafka bridges it: pause/resume The `KafkaConsumer` offers two flow-control methods: - `pause(Collection<TopicPartition>)` — subsequent `poll()` calls return **no records** for those partitions, but the consumer **keeps polling** (so it still heartbeats and stays in the group). - `resume(...)` — fetching restarts. Reactor Kafka's receive operator works like this: 1. It polls and feeds records into an internal queue that downstream operators drain according to their `request(n)` demand. 2. When downstream demand is satisfied and the queue has buffered enough (back-pressure threshold reached), it calls **consumer.pause()** on the assigned partitions — fetching stops, memory stays bounded. 3. As the downstream subscriber requests more and the queue drains, it calls **consumer.resume()** and polling delivers records again. Thus a slow `concatMap`/`flatMap` processing step automatically throttles broker ingestion — no manual code needed. ## Why heartbeats keep working Crucially, pause() does NOT stop the poll loop; it only stops *fetching records*. The consumer still calls poll() internally, which drives heartbeats and group membership. This is the safe way to slow down — contrast with simply *not calling poll()* in a plain client, which would blow `max.poll.interval.ms` and trigger a rebalance. ## Failure modes if backpressure is mishandled - **OOM:** if you break the demand chain (e.g. an operator that requests Long.MAX_VALUE unboundedly and buffers), records accumulate without bound. - **Rebalance storms / 'kicked from group':** if processing per poll exceeds `max.poll.interval.ms`, the broker assumes the consumer is dead and reassigns its partitions. Reactor Kafka mitigates by pausing rather than stalling the loop, but slow processing of *already-fetched* records can still exceed the interval if not offloaded. ## Tuning levers - `max.poll.records` — fewer records per poll = finer-grained, lower memory. - `ReceiverOptions` poll timeout and commit settings. - Reactor `prefetch`/`limitRate(n)` to cap demand. - Bounded `flatMap(..., concurrency)` so in-flight work is capped. ## Vert.x parallel The Vert.x Kafka consumer exposes the same idea explicitly via `consumer.pause()`/`consumer.resume()` and a `fetch(amount)` demand API on its read stream — same underlying mechanism, different surface.

  • Why does consumer.pause() keep the consumer in the group while a simple 'don't poll' would not?
    pause() stops record fetching but the poll loop still runs, so heartbeats continue and group membership is maintained. Not polling at all silences heartbeats/processing, eventually breaching max.poll.interval.ms and triggering a rebalance.
  • What config most directly bounds memory and pause/resume granularity?
    max.poll.records caps records returned per poll, so smaller values reduce buffered records and give finer-grained flow control; combined with Reactor's limitRate/prefetch it bounds in-flight work.

Backpressure is a faucet with a float valve: when the bucket (buffer) is full, the float shuts the valve (pause); as you scoop water out (process), the float drops and the valve reopens (resume) — the tank never overflows.

saying these in an interview costs you the question

  • Claiming Kafka has a native per-record pull/backpressure protocol (it's batch-push; pause/resume bridges it).
  • Saying pause() stops the poll loop entirely (it stops fetching but heartbeats continue).
  • Ignoring max.poll.interval.ms and assuming slow processing is always safe under backpressure.
  • Believing backpressure is automatic even if you request Long.MAX_VALUE / break the demand chain.

context