Kafka courseLesson 5 of 10
Kafka course · Lesson 5 of 10
Kafka Consumer Groups and Rebalancing
How the group coordinator manages Kafka consumer groups, eager versus cooperative versus KIP-848 rebalancing, the cooperative sticky assignor, static membership and partition moves.
On this page
- Sample setup
- Group coordinator
- How the coordinator is chosen
- Group states
- Pitfalls
- In interviews
- Rebalancing protocols
- Classic, eager
- Classic, cooperative
- Consumer protocol (KIP-848)
- Choosing and migrating
- Pitfalls
- In interviews
- Cooperative sticky assignor
- How it works
- Migrating from an eager assignor
- Pitfalls
- In interviews
- Static membership
- Worked example
- Trade-off
- Pitfalls
- In interviews
- Partition reassignment
- Moving ownership between consumers safely
- Moving replicas between brokers
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
A consumer group is only as good as its rebalances. Every time a consumer joins, leaves, crashes or restarts, partitions move, and while they move nobody processes them. In a pipeline with dozens of consumers and frequent deployments, rebalancing decides whether you see brief blips or minutes of stalled processing and duplicate work. This lesson explains who runs the group, the three rebalance protocols Kafka 4.x offers, and the settings that keep partitions where they are.
Sample setup
The experiments use the six-partition orders topic from earlier lessons on a three-broker Kafka 4.3.1 cluster. A small Java program starts consumers A, B and C in one group, eight seconds apart, then stops B, and logs every callback of its ConsumerRebalanceListener:
c.subscribe(List.of("orders"), new ConsumerRebalanceListener() {
public void onPartitionsRevoked(Collection<TopicPartition> ps) {
if (!ps.isEmpty()) log(who, "revoked " + names(ps));
}
public void onPartitionsAssigned(Collection<TopicPartition> ps) {
log(who, "assigned " + names(ps) + " -> owns " + names(c.assignment()));
}
});
while (!stop.get()) c.poll(Duration.ofMillis(100));
The same program ran three times, changing only the protocol or assignor: partition.assignment.strategy=RangeAssignor (eager), CooperativeStickyAssignor (cooperative), and group.protocol=consumer (KIP-848).
Group coordinator
What it is. Each consumer group is managed by one broker, its group coordinator. The coordinator tracks members and their liveness, runs rebalances, and stores the group’s committed offsets and metadata.
How the coordinator is chosen
Group state lives in the internal topic __consumer_offsets, which has 50 partitions by default (offsets.topic.num.partitions). A group is mapped to one partition by hashing its group.id, and the leader of that partition is the coordinator. The mapping uses Java’s String.hashCode():
def java_string_hash(s: str) -> int:
h = 0
for ch in s:
h = (31 * h + ord(ch)) & 0xFFFFFFFF
return h
def offsets_partition(group_id: str, num_partitions: int = 50) -> int:
return (java_string_hash(group_id) & 0x7FFFFFFF) % num_partitions
for group in ["payments-loader", "payments-audit", "analytics"]:
print(f"{group:16} -> __consumer_offsets partition {offsets_partition(group)}")
payments-loader -> __consumer_offsets partition 5
payments-audit -> __consumer_offsets partition 11
analytics -> __consumer_offsets partition 38
On the cluster, partitions 5, 11 and 38 were all led by broker 2 after a restart, and the tool reported broker 2 as coordinator for all three groups:
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:19092 --describe --group payments-loader --state
GROUP COORDINATOR (ID) ASSIGNMENT-STRATEGY STATE #MEMBERS
payments-loader localhost:29092 (2) uniform Empty 0
Earlier, before the restart, the same group’s coordinator was broker 3. When a coordinator’s broker fails, leadership of its __consumer_offsets partition moves, the new leader loads the group state, and clients find it again (you may briefly see COORDINATOR_LOAD_IN_PROGRESS or NOT_COORDINATOR errors, which clients retry).
Group states
| State | Meaning |
|---|---|
Empty |
no members; offsets kept until offsets.retention.minutes |
Stable |
members have their assignment and are consuming |
PreparingRebalance / CompletingRebalance |
classic protocol rebalance in progress |
Assigning / Reconciling |
consumer protocol: new target assignment being computed or delivered |
Dead |
group removed |
Pitfalls
- Hot coordinators. Many busy groups can hash to partitions led by the same broker. After broker restarts, leaders may stay concentrated until automatic preferred-leader rebalancing runs (
auto.leader.rebalance.enable, checked every 300 seconds by default). - Under-replicated
__consumer_offsets. If its replication factor is 1 (fine on a laptop, not in production) and that broker dies, groups lose their offsets. Keepoffsets.topic.replication.factorat 3.
In interviews
“Where are consumer offsets stored and who manages the group?” Answer: in __consumer_offsets; the leader of the partition that the group ID hashes to is the coordinator; it handles membership, rebalancing and offset commits.
Rebalancing protocols
What it is. A rebalance recomputes which member owns which partition. Kafka 4.x has two group protocols, and the classic one has two styles.
| Classic, eager | Classic, cooperative | Consumer protocol (KIP-848) | |
|---|---|---|---|
| Enabled by | partition.assignment.strategy = Range, RoundRobin or Sticky |
CooperativeStickyAssignor |
group.protocol=consumer |
| Who assigns | group leader (a consumer) | group leader (a consumer) | the coordinator (broker) |
| Partitions paused | all, for every member | only those that move | only those that move |
| Synchronisation | global barrier: JoinGroup then SyncGroup | two rounds of JoinGroup/SyncGroup | none; each member reconciles via heartbeats |
| Session and heartbeat settings | client: session.timeout.ms (45 s), heartbeat.interval.ms (3 s) |
same | broker: group.consumer.session.timeout.ms (45 s), group.consumer.heartbeat.interval.ms (5 s) |
| Status in 4.3 | supported, client default | supported, client default list includes it | GA since 4.0, opt-in on the client |
Classic, eager
Every member stops consuming and revokes all its partitions, rejoins, waits for the leader to compute a new assignment, then resumes. The real timeline when C joined a group of A and B:
16.0s -- C joins
17.6s B revoked [p3,p4,p5]
17.6s A revoked [p0,p1,p2]
17.7s C assigned [p4,p5] -> owns [p4,p5]
17.7s A assigned [p0,p1] -> owns [p0,p1]
17.7s B assigned [p2,p3] -> owns [p2,p3]
All six partitions stopped, although only two needed to move (and range assignment moved p2 and p3 to new owners as well). When B left, A and C again revoked everything before getting three partitions each.
Classic, cooperative
Members keep partitions that are not moving. A first rebalance revokes only the partitions that must move; a second hands them to their new owner:
16.0s -- C joins
19.0s C assigned [] -> owns []
19.0s A revoked [p2]
19.0s A assigned [] -> owns [p0,p1]
19.0s B revoked [p5]
19.0s B assigned [] -> owns [p3,p4]
22.0s C assigned [p2,p5] -> owns [p2,p5]
24.0s -- B leaves
24.1s B revoked [p3,p4]
25.0s C assigned [p4] -> owns [p2,p4,p5]
25.0s A assigned [p3] -> owns [p0,p1,p3]
Only p2 and p5 paused when C joined, and only B’s partitions moved when it left. A kept p0 and p1 throughout.
Consumer protocol (KIP-848)
The coordinator computes a target assignment and each member converges to it through its regular heartbeats. There is no group-wide barrier, and slow members do not hold up others:
16.0s -- C joins
16.2s C assigned [] -> owns []
18.4s B revoked [p5]
18.4s B assigned [] -> owns [p3,p4]
21.0s A revoked [p2]
21.0s A assigned [] -> owns [p0,p1]
21.3s C assigned [p2,p5] -> owns [p2,p5]
24.0s -- B leaves
24.1s B revoked [p3,p4]
26.3s A assigned [p4] -> owns [p0,p1,p4]
26.4s C assigned [p3] -> owns [p2,p3,p5]
The movement is incremental like cooperative rebalancing, but B and A released their partitions at different times, each on its own heartbeat. A partition is given to its new owner only after the old owner has revoked it.
Choosing and migrating
- The server enables protocols with
group.coordinator.rebalance.protocols(defaultclassic,consumer,streamsin 4.3). The Java client still defaults togroup.protocol=classicin 4.3; the documented plan is for the consumer protocol to become the client default in Kafka 5.0 and the only client option in 6.0. - With
group.protocol=consumer, client settingspartition.assignment.strategy,session.timeout.msandheartbeat.interval.msare not used. Pick a server-side assignor withgroup.remote.assignor(uniformby default, orrange). - A group can be upgraded online by rolling consumers with
group.protocol=consumer, as long as the classic group does not use an assignor with custom metadata. Client-side custom assignors are not supported by the new protocol. - Kafka Streams has its own rebalance protocol (KIP-1071, the
streamsvalue above); Streams applications useclassicunless configured otherwise.
Pitfalls
- Rebalance storms from consumers exceeding
max.poll.interval.ms: every eviction triggers another rebalance. - Assuming the new protocol fixes slow processing. It makes rebalances cheaper; it does not stop evictions caused by long processing.
- Using a librdkafka client and expecting identical behaviour. Non-Java clients add protocol support on their own release schedule; check your client’s documentation.
In interviews
Expect “what happens during a rebalance?” Strong answers distinguish eager stop-the-world rebalancing, cooperative incremental rebalancing, and the broker-driven KIP-848 protocol, and name its status: GA in 4.0, opt-in on the Java client in 4.x.
Cooperative sticky assignor
What it is. CooperativeStickyAssignor is the classic-protocol assignor that combines balance (partition counts differ by at most one), stickiness (members keep what they had) and the cooperative two-phase rebalance shown above.
How it works
- On a membership change, every member keeps consuming and sends its current ownership with its JoinGroup request.
- The leader computes a balanced target that keeps as many existing assignments as possible.
- Partitions that must move are first only revoked from their current owners (they are not yet assigned anywhere).
- The revoking members immediately rejoin, triggering a second rebalance in which the freed partitions are assigned to their new owners.
Partitions that stay put are never paused. In the cooperative run above, the join of C paused two partitions out of six; with eager range assignment, all six paused.
Migrating from an eager assignor
Members using eager and cooperative assignors cannot be mixed directly, so the upgrade is two rolling restarts:
- Roll every consumer with
partition.assignment.strategy=[CooperativeStickyAssignor, RangeAssignor](both listed). The group keeps using range, which all members still support. - Roll again with only
CooperativeStickyAssignor. Once every member supports only the cooperative assignor, the group switches.
The Java default list, [RangeAssignor, CooperativeStickyAssignor], exists to make this path easy; a group of members with the default list uses range.
Pitfalls
- Rebalance listeners written for eager mode. In cooperative mode
onPartitionsRevokedreceives only the partitions being taken away, not all of them. Code that clears all state on revocation will discard state for partitions it still owns. - Two rebalances per change. Each change involves two short rounds; this is normal.
- Equal counts are not equal load. Like every built-in assignor, it balances partition counts, not traffic.
In interviews
Explain stickiness and the two-phase revoke-then-assign mechanism, and why it beats eager rebalancing for stateful consumers. Bonus: the two-step rolling upgrade.
Static membership
What it is. By default a consumer gets a new member ID every time it starts, so a restart looks like “member left, new member joined” and triggers two rebalances. With static membership, each instance sets a stable group.instance.id (for example the pod name orders-loader-0). A restart within the session timeout is recognised as the same member and gets its previous partitions back without a rebalance.
Worked example
Two workers shared the six partitions, then worker-2 stopped for ten seconds and started again. Both runs used the consumer protocol; the only difference was group.instance.id.
Dynamic members:
20.0s -- worker-2 restarts (down for 10 s)
20.1s worker-2 revoked [p3,p4,p5]
20.4s worker-1 assigned [p3,p4,p5] -> owns [p0,p1,p2,p3,p4,p5]
30.5s worker-1 revoked [p3,p4,p5]
35.5s worker-2 assigned [p3,p4,p5] -> owns [p3,p4,p5]
Static members (group.instance.id=worker-1 and worker-2; in this run worker-2 happened to own p0 to p2):
20.1s -- worker-2 restarts (down for 10 s)
20.1s worker-2 revoked [p0,p1,p2]
30.6s worker-2 assigned [p0,p1,p2] -> owns [p0,p1,p2]
With dynamic membership, the partitions moved to worker-1 and back: two rebalances, and work done twice for stateful consumers (caches rebuilt, state reloaded). With static membership, worker-1 was not disturbed, and worker-2 got exactly its own partitions back. The revoked line is the client’s local callback on shutdown; the coordinator kept the assignment reserved for the instance ID.
Trade-off
During the ten seconds, nobody consumed p0 to p2, so their lag grew. Static membership trades a short pause on one member’s partitions for no reshuffle. If the instance does not return within the session timeout (45 seconds by default under both protocols; session.timeout.ms for classic, group.consumer.session.timeout.ms on the broker for the consumer protocol), it is removed and its partitions are reassigned.
Pitfalls
- Duplicate instance IDs. Two running consumers with the same
group.instance.idcannot both be members: under the classic protocol one is fenced withFencedInstanceIdException, and the consumer protocol rejects the second one withUnreleasedInstanceIdExceptionwhile the first is still active. Derive the ID from something unique and stable, such as a StatefulSet pod name. - Very long session timeouts delay real failover: a crashed static member’s partitions sit idle until the timeout.
max.poll.interval.msstill applies. A static member that stops polling stops heartbeating, and its partitions move after the session timeout.
In interviews
“How do you stop rolling deployments from causing rebalance storms?” Static membership with a stable group.instance.id and a session timeout longer than a restart, plus cooperative or KIP-848 rebalancing so that the rebalances that do happen are incremental.
Partition reassignment
What it is. Partitions get reassigned in two different senses, and interviewers use both:
- Between consumers: a rebalance moves ownership of a partition from one group member to another.
- Between brokers: an operator moves a partition’s replicas to different brokers, for example to use a new broker or drain an old one.
Moving ownership between consumers safely
When a partition is taken from consumer A and given to consumer B, B starts from the group’s last committed offset. Everything A processed but did not commit is processed again by B. To keep that window small:
| Callback | What to do |
|---|---|
onPartitionsRevoked(partitions) |
finish or abandon in-flight work for these partitions, then commitSync() their offsets and flush their state |
onPartitionsAssigned(partitions) |
load state, or seek() to an externally stored offset if you keep offsets in your sink |
onPartitionsLost(partitions) |
the partitions were already given to someone else (for example after a session timeout): do not commit, just discard local state |
Adding partitions to a topic also causes reassignment: consumers notice the new partitions on their next metadata refresh (metadata.max.age.ms, 5 minutes by default) and a rebalance assigns them. New partitions with no committed offsets start from auto.offset.reset; with the default latest, records produced to them before the group noticed can be skipped.
Moving replicas between brokers
kafka-reassign-partitions.sh takes a JSON plan, copies data to the new replicas, and switches over once they are in sync. Moving keytest partition 0 from broker 3 to broker 1:
cat > move.json <<'EOF'
{"version":1,"partitions":[{"topic":"keytest","partition":0,"replicas":[1]}]}
EOF
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:19092 --reassignment-json-file move.json --execute
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:19092 --reassignment-json-file move.json --verify
Current partition replica assignment
{"version":1,"partitions":[{"topic":"keytest","partition":0,"replicas":[3],"log_dirs":[".../cluster/data3"]}]}
Save this to use as the --reassignment-json-file option during rollback
Successfully started partition reassignment for keytest-0
Status of partition reassignment:
Reassignment of partition keytest-0 is completed.
Clearing broker-level throttles on brokers 1,2,3
Clearing topic-level throttles on topic keytest
Afterwards --describe showed Leader: 1 Replicas: 1 Isr: 1. (The log directory path is shortened here.) On a busy cluster, add --throttle <bytes/sec> so the copy does not starve production traffic, and keep the printed rollback plan. The --generate option proposes a plan for a list of topics and brokers, but large clusters usually use a balancer such as Cruise Control (see operating Kafka).
Pitfalls
- Committing in
onPartitionsLost: another member may already own the partition, and your commit can rewind or overwrite its progress. - Reassigning without throttling can saturate disks and networks and push followers out of the ISR.
- Forgetting that replica moves change leaders: clients follow automatically, but leadership can end up uneven; run a preferred leader election (
kafka-leader-election.sh --election-type PREFERRED) if automatic rebalancing is disabled.
In interviews
Be explicit about which kind of reassignment you mean. For consumers, talk about commit-on-revoke and the reprocessing window. For replicas, describe the reassignment tool, throttling and rollback, and tools that automate it.
Practice questions
Every deployment of your 20-instance consumer group causes several minutes of stalled processing. What do you change?
Use static membership (group.instance.id per instance, with a session timeout longer than one instance’s restart) so that restarts do not trigger rebalances, and switch to cooperative rebalancing or the KIP-848 consumer protocol so that any rebalance only pauses partitions that move. Also check that shutdown calls close() and that max.poll.interval.ms is not being exceeded during start-up.
Why does eager rebalancing pause every partition, even when only one consumer joins?
In the eager classic protocol, every member must revoke all its partitions before rejoining, because the leader computes a fresh assignment without knowing which partitions are safe to keep. Nobody consumes until the new assignment is distributed in SyncGroup. Cooperative rebalancing and the consumer protocol avoid this by revoking only partitions that change owner.
How is a consumer group’s coordinator chosen, and what happens if that broker fails?
The group ID is hashed to one of the __consumer_offsets partitions (50 by default), and the leader of that partition is the coordinator. If the broker fails, leadership moves to another replica, which loads the group’s state from the log; clients rediscover the coordinator and retry, with a short interruption.
What is the difference between onPartitionsRevoked and onPartitionsLost?
onPartitionsRevoked is called while you still own the partitions, so you can commit offsets and flush state before they move. onPartitionsLost is called when ownership has already gone, for example after the consumer’s session expired; another member may already be processing those partitions, so you must not commit, only clean up.
A team wants to move to group.protocol=consumer. What should they check?
That brokers are on Kafka 4.0 or later with the consumer protocol enabled, that their clients support it, that they do not rely on a custom client-side assignor (not supported), that code does not set the now-ignored session.timeout.ms, heartbeat.interval.ms or partition.assignment.strategy, and that rack-aware assignment is not required, since it is not fully supported yet. The group can then be migrated by a rolling restart.
What does static membership trade away?
Faster failover. While a static member is down, its partitions are not reassigned until the session timeout expires, so their lag grows. In exchange, restarts within that window cause no rebalance and the member keeps its partitions and any local state.
Key takeaways
- The group coordinator is the leader of the
__consumer_offsetspartition that the group ID hashes to; it manages membership, rebalances and offsets. - Eager rebalancing pauses every partition; cooperative sticky and the KIP-848 consumer protocol pause only partitions that move.
- KIP-848 has been GA since Kafka 4.0, moves assignment to the broker and removes the global barrier, but the Java client still defaults to
classicin 4.x. - Migrate to
CooperativeStickyAssignorwith two rolling restarts, and remember that revocation callbacks then receive only moving partitions. - Static membership (
group.instance.id) lets restarts within the session timeout skip rebalancing, at the cost of slower failover. - Commit in
onPartitionsRevoked, never inonPartitionsLost; move replicas between brokers with throttledkafka-reassign-partitions.shor a balancer.
Progress is saved in this browser only. No account needed.

