How does a ConsumerRebalanceListener's callback behavior differ between eager and cooperative rebalancing, and where should you commit offsets?
answer
- revoked = commit here (commitSync)
- assigned = init/seek here
- lost = do NOT commit (fenced)
- eager: full set every time; cooperative: only the delta
- auto-commit fires before revoked
basics
~10 sConsumerRebalanceListener has onPartitionsRevoked, onPartitionsAssigned, and onPartitionsLost. Commit offsets in onPartitionsRevoked before giving partitions up. Under eager you're given all partitions; under cooperative you only get the ones actually moving.
solid answer
~40 sYou register a ConsumerRebalanceListener when subscribing to react to assignment changes. onPartitionsRevoked fires just before partitions are taken away — this is the place to commit offsets and flush state so the next owner resumes cleanly. onPartitionsAssigned fires when you receive partitions and is where you'd initialize state or seek. onPartitionsLost (KIP-429/KIP-345) fires when partitions were taken without a clean revoke (e.g. the member was fenced), so you should NOT commit there because you no longer own them. The protocol changes the SET passed: under eager, onPartitionsRevoked gets ALL currently owned partitions every rebalance and onPartitionsAssigned gets the full new set; under cooperative, onPartitionsRevoked gets only the partitions actually being moved away and onPartitionsAssigned gets only newly added ones, so callbacks are smaller and cheaper.
go deeper
Know the three callbacks exist and that you commit offsets in onPartitionsRevoked.
Explain assigned vs revoked vs lost and why lost must not commit.
Articulate how eager passes the full set while cooperative passes only the delta, and write subset-safe revoke logic.
Design exactly-once-aware handover (external offset stores, transactional flush) across protocols and reason about callback latency in the poll loop.
## The listener `ConsumerRebalanceListener` is an interface you pass to `consumer.subscribe(topics, listener)`. It has three callbacks the consumer invokes during rebalances: - **`onPartitionsRevoked(partitions)`** — called *before* the consumer loses ownership of `partitions`. This is the canonical place to **commit offsets** (and flush any local state / external buffers) so whoever picks up the partition continues from the right point and you avoid reprocessing. - **`onPartitionsAssigned(partitions)`** — called *after* the consumer gains `partitions`, before the next `poll()` returns records for them. Use it to initialize per-partition state, restore caches, or `seek()` to a custom offset. - **`onPartitionsLost(partitions)`** — called when partitions were lost **without** a graceful revoke — e.g. the member was fenced out (a newer generation took over) or it exceeded `max.poll.interval.ms`. The default implementation calls `onPartitionsRevoked`, but you typically override it to **NOT commit**, since you no longer safely own those partitions and committing could clobber the new owner's progress. ## How the protocol changes the callbacks ### Eager Every rebalance, before anything else, the consumer revokes **everything**. So `onPartitionsRevoked` receives the consumer's **entire** current assignment each time, and `onPartitionsAssigned` later receives the **entire** new assignment — even partitions it already had and kept. Callbacks are 'big' and fire on every rebalance for all partitions. ### Cooperative The consumer keeps the partitions it retains. So `onPartitionsRevoked` receives **only the partitions actually being moved away** (often empty), and `onPartitionsAssigned` receives **only the newly added** partitions. Retained partitions don't appear in either callback and never pause. This is what makes cooperative rebalances cheap — and it means your revoke logic should be written to handle a *subset*, not assume it's the whole assignment. ## Where to commit - **Auto-commit (`enable.auto.commit=true`):** the consumer auto-commits before `onPartitionsRevoked` for you. Fine for at-least-once with simple processing. - **Manual commit:** call `consumer.commitSync(offsets)` inside `onPartitionsRevoked`. Use `commitSync` (not async) here because you're about to lose the partitions and want the commit to complete first. - **Never commit in `onPartitionsLost`** for the lost partitions — ownership is already gone. ## Edge cases - In cooperative mode, since revoke is incremental, a long-running, blocking commit in `onPartitionsRevoked` still blocks the poll loop — keep it tight. - `onPartitionsAssigned` is also where you should re-seek if you store offsets externally (e.g. in a DB) rather than in Kafka. - If you do external transactional processing, the revoke callback is your last chance to flush atomically.
- Why should you not commit offsets in onPartitionsLost?Because onPartitionsLost means the partitions were taken from you without a clean revoke — typically you were fenced and a new generation already owns them. Committing would write offsets for partitions you no longer own, potentially overwriting the new owner's progress and causing skipped or reprocessed records.
- In cooperative mode, why might onPartitionsRevoked be called with an empty set?Because cooperative rebalancing only revokes partitions that actually need to move. If a particular consumer retains all its partitions across the rebalance, none are being moved away from it, so its onPartitionsRevoked receives an empty collection.
saying these in an interview costs you the question
- Committing offsets in onPartitionsAssigned instead of onPartitionsRevoked.
- Assuming onPartitionsRevoked always receives the full assignment (true only for eager).
- Committing in onPartitionsLost (you no longer own those partitions).
- Using commitAsync in the revoke callback and not waiting — you may lose the commit before partitions move.