Explain the pause()/resume() backpressure pattern for offloading long-running record processing to worker threads while keeping the consumer in its group.
answer
- poll thread polls, workers process
- pause() = poll alive but no new records
- high/low watermark backpressure
- commit highest CONTIGUOUS completed offset
- rebalance clears pause; RebalanceListener to drain/rewind
- ordering lost unless key-sharded
basics
~20 sHand records to a worker pool for slow processing. To avoid fetching more than you can handle, call consumer.pause() on the assigned partitions and keep calling poll() (which now returns nothing but proves liveness). When workers catch up, call resume(). Commit only completed offsets.
solid answer
~50 sThe pattern decouples fetch liveness from processing. The poll thread keeps calling poll() on the required cadence (so max.poll.interval.ms is never breached), but actual processing runs on a separate executor. To prevent unbounded buffering, you apply backpressure: when in-flight work exceeds a threshold, call consumer.pause(assignment) — poll() then returns no new records for those partitions but still sends heartbeats and refreshes the poll deadline. When workers drain, call resume(). You must commit only offsets that have actually completed (track the highest contiguous committed offset per partition), and on a rebalance you have to pause/seek back to the last committed offset for revoked partitions to avoid skipping un-processed records. This keeps the consumer healthy under long processing but introduces complexity: ordering within a partition can be lost if multiple workers process it concurrently, offset tracking must be careful, and you need a ConsumerRebalanceListener to drain or rewind in-flight work. Frameworks like Spring Kafka's async/ack modes or the Confluent Parallel Consumer implement this for you.
go deeper
Know the gist: do slow work on other threads and keep calling poll() so you don't get kicked out.
Explain pause()/resume() as backpressure (poll stays alive, no new records) and that you commit only completed offsets.
Detail the high/low watermark loop, contiguous-offset commit, rebalance handling via ConsumerRebalanceListener, and the ordering tradeoff.
Weigh this against simpler configs, design the offset-tracking and rebalance-drain protocol, and decide build-vs-adopt (Parallel Consumer / Spring Kafka async) with ordering/throughput SLAs.
## The core idea: keep poll() alive, move work off the loop The single-threaded poll loop ties processing time to `max.poll.interval.ms`. The fix for genuinely long work is to make the poll thread do almost nothing but poll, and process elsewhere: ```java while (running) { ConsumerRecords<K,V> records = consumer.poll(Duration.ofMillis(200)); for (var r : records) workerPool.submit(() -> process(r)); // offload // backpressure: if (inFlight() > HIGH_WATERMARK) consumer.pause(consumer.assignment()); else if (inFlight() < LOW_WATERMARK) consumer.resume(consumer.assignment()); commitCompletedOffsets(); // only what workers finished } ``` Now `poll()` is called every ~200 ms regardless of how slow processing is, so `max.poll.interval.ms` is never threatened. ## Why pause()/resume() are essential If you keep calling `poll()` without pausing, the consumer keeps fetching new records into the application — but your workers can't keep up, so an unbounded queue grows and you OOM. `consumer.pause(partitions)` tells the consumer: keep the loop alive (heartbeats, poll-deadline refresh, metadata) but **return no new records** for those partitions. When workers drain below a low watermark, `consumer.resume(partitions)` re-enables fetching. This is classic high/low-watermark backpressure. Important: a paused partition is still *assigned* — pausing is not unsubscribing. A rebalance clears the paused state, so you must re-apply pause after partitions are reassigned. ## Offset management is the hard part With async workers, records complete out of order. You must commit only the **highest contiguous completed offset** per partition — committing a higher offset whose predecessor hasn't finished would lose data on a crash. Maintain a per-partition structure (e.g. a sorted set of completed offsets) and advance the commit point only across the contiguous prefix. ## Rebalance correctness Use a `ConsumerRebalanceListener`: - `onPartitionsRevoked`: stop accepting new work for revoked partitions, optionally wait for in-flight work to finish, and commit completed offsets — otherwise the new owner reprocesses (acceptable for at-least-once) or, worse, you commit ahead of incomplete work and lose records. - `onPartitionsAssigned`: re-apply any needed pause state; seek if you track offsets externally. ## Ordering tradeoff Kafka guarantees order **within a partition**. If you fan out a single partition's records to multiple threads, you lose that ordering unless you key-shard work to a stable worker (hash(key) → worker). The Confluent Parallel Consumer offers KEY / PARTITION / UNORDERED ordering modes precisely to manage this. ## When to use what - Bounded long processing, order matters, low complexity → raise `max.poll.interval.ms` and/or lower `max.poll.records`. - High throughput with slow per-record work, can tolerate keyed parallelism → pause/resume + worker pool, or adopt the Parallel Consumer / Spring Kafka async ack. ## Pitfalls - Forgetting to re-pause after a rebalance → buffer blow-up. - Committing offsets for records still in flight → silent data loss on crash. - Doing blocking work on the poll thread anyway (e.g. a slow `commitSync` or a synchronous side call) → defeats the whole point.
- If you keep polling but never pause, what goes wrong?poll() keeps returning new records faster than your workers process them, so an in-memory queue grows unbounded and the JVM eventually OOMs. pause() applies backpressure so fetching stops while workers catch up.
- How do you commit offsets safely when workers complete out of order?Track completed offsets per partition and commit only the highest *contiguous* completed offset (the committed point is the largest N such that all offsets up to N are done). Committing a non-contiguous higher offset would lose the gap's records if the consumer crashes.
- What does a rebalance do to paused partitions and in-flight work?A rebalance clears paused state and may revoke partitions mid-flight. You need a ConsumerRebalanceListener: in onPartitionsRevoked drain/commit completed work, and in onPartitionsAssigned re-apply pause. Otherwise you over-buffer or lose/duplicate records.
saying these in an interview costs you the question
- Saying you can just spawn threads and ignore pause() — you'll OOM from unbounded buffering.
- Committing the offset of every record as workers finish, ignoring contiguity — causes data loss on crash.
- Assuming pause() leaves the consumer healthy without re-applying it after a rebalance — rebalance clears pause state.
- Claiming per-partition ordering is preserved when fanning one partition across multiple threads — it isn't unless you key-shard.