skip to content

When you decouple the poll loop from a worker thread pool, what offset-commit hazard arises and how do you avoid losing or reprocessing records?

level: seniorimportance: must knowfreq 55%

answer

  1. offset = high-water mark, not per-record ack
  2. commit highest CONTIGUOUS completed offset
  3. out-of-order completion -> skip = data loss
  4. pause()/resume() for backpressure
  5. route same key to same worker for order

basics

~20 s

Workers finish out of offset order, so committing the latest offset can skip records still in flight on slower workers. Avoid it by committing only the highest contiguous completed offset per partition, and pause/resume the partition to bound in-flight work.

solid answer

~50 s

If you poll on one thread and dispatch records to a worker pool, work completes out of order — worker B may finish offset 105 while worker A is still on offset 103. Kafka commits an offset meaning 'everything below this is done', so committing 106 marks 103 as processed; a crash then loses 103 (data loss). Naively committing only the lowest in-flight offset risks reprocessing. The safe rule: track per-partition completion and commit only the highest *contiguous* offset whose predecessors all completed. You also lose the consumer's natural in-order, at-most-one-poll-of-inflight backpressure, so you must pause(partition) when the in-flight buffer is full and resume() when it drains, otherwise poll() keeps fetching and you can exceed memory or max.poll.interval.ms. Per-key ordering is also lost unless you route records of the same key to the same worker. This is exactly the complexity the Confluent Parallel Consumer packages.

go deeper

for a junior

Understand that offsets are a single 'everything-below-is-done' marker, not per-record acks.

for a middle

Explain why out-of-order completion plus naive commit causes data loss.

for a senior

Design the contiguous-prefix commit plus pause/resume backpressure and key routing.

for a principal

Weigh building this vs adopting Confluent Parallel Consumer, and integrate with rebalance/idempotency guarantees end to end.

## Why decouple at all The one-consumer-per-thread model caps parallelism at the partition count and processes each partition serially. If processing is slow (e.g. an I/O call per record) and you can't add partitions, you separate concerns: a **poll thread** owns the (non-thread-safe) consumer and only does poll/commit/pause/resume; a **worker pool** does the slow processing. This breaks the partition ceiling — many workers can process one partition's records concurrently. ## The core hazard: offsets are a high-water mark A Kafka committed offset is a single number per partition meaning **'all records with offset < this are processed.'** It is NOT a set of individual acks. With a worker pool, completion order != offset order: - poll returns offsets 100,101,102,103 for partition P. - Workers run them concurrently. Worker for 103 finishes first. - If you commit 104 now, you've told Kafka 100-103 are done. But 100-102 may still be running. - Crash here -> on restart you resume at 104 -> **100-102 are lost** (never reprocessed). That's silent data loss. The opposite mistake — always committing the smallest in-flight offset — causes large reprocessing windows and stalls progress behind one slow record. ## The safe commit rule: highest contiguous completed offset Maintain, per partition, the set of completed offsets. Commit the **highest offset N such that every offset below N is completed** (the contiguous prefix). If 100,101,103 are done but 102 isn't, you may only commit 102 (i.e. 'through 101'). When 102 completes you can jump to 104. This guarantees at-least-once with no gap-skipping. (Combine with idempotent processing to make reprocessing harmless.) ## Backpressure: pause/resume The single-thread loop gives natural backpressure — you don't poll again until you've processed. A worker pool removes that: `poll()` keeps returning records even if workers are saturated, growing an unbounded in-flight buffer and risking OOM, and long gaps between *useful* polls can still trip `max.poll.interval.ms`. Fix: when in-flight count for a partition exceeds a threshold, call `consumer.pause(partitions)`; when it drains, `consumer.resume(partitions)`. You keep calling poll() (to stay in the group and send heartbeats) but paused partitions return nothing. ## Ordering is lost unless you route by key Concurrent workers destroy per-key/per-partition order. If order matters, hash the record key to a fixed worker (or queue) so all records of a key are processed sequentially while different keys run in parallel. This is exactly the 'key ordering' guarantee the Confluent Parallel Consumer provides. ## Rebalance interaction On a partition revocation you must wait for or cancel in-flight work for that partition and commit the safe offset in `onPartitionsRevoked`, or another instance will reprocess from the last commit. Cooperative rebalancing limits the blast radius. ## Summary checklist 1. Commit only the contiguous completed prefix per partition. 2. Pause/resume to bound in-flight records. 3. Route by key if order matters. 4. Make processing idempotent (at-least-once). 5. Drain/commit on revocation.

  • Why can't you just commit the offset of whichever record finishes last?
    Because the committed offset means everything below it is done. If a later offset finishes first and you commit it, you mark still-in-flight earlier offsets as processed, losing them on a crash.
  • Your in-flight queue grows unbounded and the consumer gets kicked from the group. What two mechanisms fix this?
    pause()/resume() to stop fetching when workers are saturated (bounding memory), and keeping max.poll.interval.ms in check so the consumer isn't evicted for slow progress.

saying these in an interview costs you the question

  • Committing the latest polled offset regardless of which workers finished.
  • Treating a Kafka offset as a per-message ack set rather than a contiguous high-water mark.
  • Forgetting backpressure, letting the in-flight buffer grow unbounded.
  • Assuming a worker pool preserves per-key ordering without explicit key routing.

context