Distributed Reliability · Staff
Kafka lag grows while consumers are alive. Why do partitions keep moving?
Take a few minutes to form your approach. Then open a worked answer and compare the decisions.
Reveal a worked answer
Look at the time between calls to poll(), not just process health. A consumer may fetch a batch of AI indexing jobs and spend six minutes parsing one huge PDF. If its max.poll.interval.ms is five minutes and it does not poll again, the group considers it unable to make progress and starts reassignment. Kafka's consumer configuration documents the poll interval and notes that the exact reassignment timing also depends on membership protocol and session timeout. A process can still be running its parser after it has lost partition ownership.
That creates two risks. Progress falls because partitions rebalance repeatedly and new owners replay uncommitted records. Also, the old worker may finish a document and try to publish an index update after a new owner has started the same work. Offsets alone do not prevent duplicate side effects outside Kafka. Track poll gaps, rebalance frequency, partition revocations, in-flight job durations and event-to-index commit positions. Separate "consumer is alive" from "consumer currently owns this partition."
Keep polling responsive by moving slow processing to bounded workers or by limiting batch size, but do not commit a record's offset before its required downstream effect is durable. If processing is asynchronous, maintain per-partition ordering and a safe contiguous commit watermark. On revocation, stop or fence work that can write stale results, or make writes idempotent and version-checked against source document revision. Raising the poll interval can reduce churn if the work is predictably slow, but it also delays detection of a truly stuck worker. Size it from measured processing tails, not one average PDF.
The interviewer may say a new consumer will simply redo the work, so no data is lost. Maybe. It can also publish an older document version after the replacement worker processed a newer one. The queue lease expires while an agent tool call is running covers a queue lease expiring during an agent tool call. Kafka's poll liveness and partition ownership are the boundary here, and the indexer must prove that only a valid document version becomes searchable.
Continue reading
Related questions
Read beyond the question
Explore more distributed reliability
Follow another question in this area, or search the complete Question Library.
Browse this area →Browse Question Library →