Explain the ConsumerRebalanceListener callbacks (onPartitionsRevoked, onPartitionsAssigned, onPartitionsLost) and what you should do in each.
answer
- Revoked = last chance, commit + flush
- Assigned = init state + seek()
- Lost = involuntary, DON'T commit
- Lost default delegates to Revoked (override it!)
- callbacks run on poll thread - keep fast
basics
~20 sIt's a callback you pass to subscribe() so you can react to rebalances: onPartitionsRevoked (commit offsets / flush state before losing partitions), onPartitionsAssigned (init state / seek for new partitions), onPartitionsLost (partitions gone unexpectedly — don't commit).
solid answer
~50 sConsumerRebalanceListener hooks into group rebalances. onPartitionsRevoked(partitions) fires before partitions are taken away — your last chance to commit offsets and flush any external state for those partitions; it runs synchronously on the poll thread. onPartitionsAssigned(partitions) fires after you gain partitions — initialize per-partition state, restore cached offsets, or call seek() to override the start position. onPartitionsLost(partitions) (added in KIP-429) fires when partitions are lost involuntarily — e.g. the consumer fell out of the group (session timeout) so another member may already own them; you must NOT commit (you'd clobber the new owner) — just clean up local state. Default onPartitionsLost delegates to onPartitionsRevoked, so override it if revoke does a commit. With eager rebalancing all partitions are revoked each time; with cooperative (incremental) rebalancing only the partitions actually moving are revoked/assigned, so these callbacks receive smaller sets.
go deeper
Recognize the three callback names and that they relate to rebalances.
Know revoked=commit-before-loss, assigned=init/seek, lost=involuntary.
Articulate the don't-commit-in-lost rule, the default-delegation gotcha, and poll-thread timing.
Design state-handoff/exactly-once flows using seek() + external offset stores and reason about cooperative vs eager callback granularity.
## ConsumerRebalanceListener When you use group-managed subscription (`subscribe()`), Kafka can **reassign partitions among members** at any time (a rebalance). The `ConsumerRebalanceListener` is the callback interface that lets your code react around those transitions. You register it via `subscribe(topics, listener)`. All callbacks run **synchronously on the thread that calls poll()** — they block the poll loop, so keep them fast. ### onPartitionsRevoked(Collection<TopicPartition>) Called **before** the consumer gives up the listed partitions (during a rebalance). This is your **last chance** while you still own them: - **Commit current offsets** for those partitions (`commitSync`) so the next owner resumes at the right place. This is the classic place to do a synchronous commit if you manage offsets manually. - **Flush/close external state** — e.g. write a buffered batch to a database, release a per-partition resource. With **eager** rebalancing (RangeAssignor/RoundRobin/StickyAssignor), *all* assigned partitions are revoked on every rebalance ("stop-the-world"). With **cooperative** rebalancing (CooperativeStickyAssignor), only the partitions that are actually being moved away are revoked. ### onPartitionsAssigned(Collection<TopicPartition>) Called **after** the rebalance completes and you have gained the listed partitions, **before** the next poll() returns records. Use it to: - **Initialize per-partition state** (load checkpoints, allocate buffers). - **seek()** to a custom starting offset — e.g. read committed offsets from your own store and `seek()` there, instead of relying on `auto.offset.reset`. ### onPartitionsLost(Collection<TopicPartition>) Introduced in **KIP-429** (cooperative rebalancing). Called when partitions are lost **involuntarily** — the consumer was kicked out of the group (e.g. it missed the `session.timeout.ms` heartbeat deadline or `max.poll.interval.ms` was exceeded), so the coordinator already reassigned those partitions to **someone else**. Critical rule: **do NOT commit offsets here.** You no longer own the partitions; another consumer may already be processing them, and committing would overwrite their progress / cause duplicate or lost work. Just **discard local state and clean up**. The default implementation of `onPartitionsLost` simply calls `onPartitionsRevoked`. That default is dangerous if your revoke does a commit — so when you commit in onPartitionsRevoked, **override onPartitionsLost** to do clean-up-only. ### Summary table | Callback | When | Do | Don't | |---|---|---|---| | onPartitionsRevoked | before giving up partitions (graceful) | commit, flush state | block for long | | onPartitionsAssigned | after gaining partitions | init state, seek() | assume offset 0 | | onPartitionsLost | partitions taken involuntarily | clean up local state | commit offsets | ### Edge cases - Callbacks run on the poll thread → long work delays heartbeats and can cause *another* rebalance. - On clean `consumer.close()`, onPartitionsRevoked fires for a final commit. - assign()-based consumers never get these callbacks — there are no rebalances.
- Why must you avoid committing offsets in onPartitionsLost?onPartitionsLost means you were ejected from the group and the partitions are likely already owned by another consumer. Committing would overwrite the new owner's offsets, causing skipped or duplicated records. You should only clean up local state there.
- Why might you call seek() inside onPartitionsAssigned?To start consuming from a position you control — e.g. offsets stored in your own database for exactly-once semantics — rather than relying on the committed offset or auto.offset.reset.
saying these in an interview costs you the question
- Committing offsets in onPartitionsLost.
- Saying onPartitionsRevoked fires after partitions are already gone — it fires before, while you still own them.
- Forgetting that the default onPartitionsLost delegates to onPartitionsRevoked.
- Doing slow/blocking work in callbacks (they run on the poll thread and delay heartbeats).
- Expecting these callbacks with assign() — they only fire under group management.