skip to content

How do you commit offsets safely around a consumer group rebalance using ConsumerRebalanceListener?

level: seniorimportance: should knowfreq 60%

answer

  1. revoked -> commitSync before handoff
  2. assigned -> seek() / init state
  3. onPartitionsLost -> cleanup, do NOT commit
  4. cooperative-sticky revokes only moved partitions
  5. pass listener to subscribe()

basics

~20 s

Register 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 s

A 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

for a junior

Know that a rebalance moves partitions between consumers and a listener lets you react.

for a middle

Explain committing in onPartitionsRevoked and seeking in onPartitionsAssigned to avoid duplicates.

for a senior

Use commitSync before handoff, handle onPartitionsLost correctly, and distinguish eager vs cooperative revoke sets.

for a principal

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)

context