How do you commit offsets safely around a consumer group rebalance using ConsumerRebalanceListener?
answer
- revoked -> commitSync before handoff
- assigned -> seek() / init state
- onPartitionsLost -> cleanup, do NOT commit
- cooperative-sticky revokes only moved partitions
- pass listener to subscribe()
basics
~20 sRegister a ConsumerRebalanceListener via subscribe(). In onPartitionsRevoked, commit offsets for partitions you're about to lose before they move to another consumer. onPartitionsAssigned runs when you gain partitions (e.g., to seek to a custom position). This prevents duplicate or lost processing across rebalances.
solid answer
~40 sA rebalance reassigns partitions among group members (member join/leave, new partitions). Without care, in-flight processed-but-uncommitted records get reprocessed by the new owner. You pass a ConsumerRebalanceListener to subscribe(). onPartitionsRevoked(partitions) fires just before you give up partitions — commit your latest processed offsets here (commitSync, so it completes before handoff). onPartitionsAssigned(partitions) fires after you receive partitions — use it to initialize state or seek() to externally-stored offsets. With cooperative-sticky rebalancing (the modern default for many setups), only the partitions actually being moved are revoked, so onPartitionsRevoked receives a smaller set, and there's also onPartitionsLost for the case where partitions are taken away without a clean revoke (e.g., session timeout). The listener is the hook for transactional/external offset stores and for bounding duplicate work.
go deeper
Know that a rebalance moves partitions between consumers and a listener lets you react.
Explain committing in onPartitionsRevoked and seeking in onPartitionsAssigned to avoid duplicates.
Use commitSync before handoff, handle onPartitionsLost correctly, and distinguish eager vs cooperative revoke sets.
Integrate the listener with an external transactional offset store for effectively-once, and reason about callback latency impact on rebalance duration and max.poll.interval.ms.
## What a rebalance is A **consumer group** spreads a topic's partitions across its member consumers. A **rebalance** is the protocol that re-divides partitions when membership or partition count changes (a consumer joins, leaves, crashes/times out, or partitions are added). During a rebalance, ownership of partitions moves between consumers. ## The problem rebalances create for offsets Suppose consumer A owns partition P and has processed up to offset 100 but only committed up to 80. A rebalance moves P to consumer B. B reads the committed offset (80) and reprocesses 80–100 — **duplicates**. Symmetrically, if your app considers work done but never commits before losing the partition, the next owner redoes it. The fix: **commit before you lose the partition.** ## ConsumerRebalanceListener You register it when subscribing: ``` consumer.subscribe(topics, listener); ``` It has three callbacks (the third added later): - **onPartitionsRevoked(partitions)** — invoked **before** the consumer relinquishes these partitions (during a normal, cooperative or eager rebalance). This is where you **commitSync** the latest offsets for those partitions so the next owner resumes correctly. Use commitSync (not async) so the commit completes before handoff. - **onPartitionsAssigned(partitions)** — invoked **after** the consumer is granted partitions, before the next poll returns their records. Use it to set up per-partition state, or to `seek()` to offsets you store externally (e.g., in a DB for exactly-once with an external sink). - **onPartitionsLost(partitions)** — invoked when partitions are taken away abruptly (e.g., the consumer fell out of the group via session timeout, or max.poll.interval.ms exceeded). Here you must **not** commit (you no longer own them; another member may already have them) — you clean up local state only. If you don't override it, it defaults to calling onPartitionsRevoked, which can cause errors. ## Eager vs cooperative protocols - **Eager** (RangeAssignor/RoundRobinAssignor): every consumer revokes **all** its partitions at the start of a rebalance ('stop-the-world'), then gets a new assignment. onPartitionsRevoked sees the full set. - **Cooperative** (CooperativeStickyAssignor): only the partitions that actually move are revoked, so onPartitionsRevoked receives just the subset being reassigned, and most partitions keep flowing — far less disruption. This requires the cooperative rebalance protocol and is the recommended default for many workloads. ## Pattern ``` class Listener implements ConsumerRebalanceListener { public void onPartitionsRevoked(Collection<TopicPartition> tps) { consumer.commitSync(currentOffsets); // flush before handoff } public void onPartitionsAssigned(Collection<TopicPartition> tps) { for (var tp : tps) consumer.seek(tp, loadFromDb(tp)); // external store } public void onPartitionsLost(Collection<TopicPartition> tps) { cleanupState(tps); // do NOT commit } } ``` ## Why this matters Rebalance-time commits bound the duplicate-processing window and are the integration point for storing offsets atomically alongside your output (the foundation of effectively-once with an external transactional sink).
- Why use commitSync rather than commitAsync inside onPartitionsRevoked?The commit must complete before the partition is handed to the new owner. commitSync blocks until acknowledged; commitAsync could return before the commit lands, letting the next owner read a stale offset and reprocess.
- What is onPartitionsLost for, and how does it differ from onPartitionsRevoked?onPartitionsLost fires when partitions are taken away without a clean revoke (e.g., the member dropped out via session/max.poll timeout). You no longer own them, so you must not commit — only clean up local state. Revoked fires during an orderly rebalance where committing is correct.
saying these in an interview costs you the question
- Committing inside onPartitionsLost (you no longer own the partitions)
- Using commitAsync in onPartitionsRevoked when a guaranteed pre-handoff commit is needed
- Believing every rebalance revokes all partitions (only eager does; cooperative revokes the moved subset)
- Putting heavy/blocking work in the listener callbacks (they run on the poll thread and delay the rebalance)