Rebalance — stop-the-world vs. cooperative

Eager rebalance revokes every partition from every consumer on any membership change, while cooperative rebalance (KIP-429) only revokes partitions that are actually moving — turning a 14-minute outage into a 1-minute brief gap on a couple of lanes.

Previously

Followers and leaders are sorted. Now the OTHER source of churn: consumers come and go, and the group has to reassign partitions. Eager rebalance is a 14-minute outage; cooperative is a 1-minute brief gap.

Scene 07

Rebalance — stop-the-world vs. cooperative

  1. Watch
  2. Try it
  3. Predict
  4. Capture
Producersteady traf…Group coordinatorprotocol = eagerPARTITIONS · 6P0→ C0P1→ C1P2→ C2P3→ C3P4→ C0P5→ C1CONSUMER GROUP A · 4 CONSUMERSdynamic membership — restart triggers rebalanceC0owns P0, P4C1owns P1, P5C2owns P2C3owns P3
What to watch for

Four consumers split six partitions. Watch the lanes — they're all green (active). At tick 8, C0 exceeds max.poll.interval.ms; the group enters PreparingRebalance and every lane goes dark for a few ticks. Watch what happens to the lanes that AREN'T moving.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Coordinator.onJoinGroup # eager
every member revokes everything before reassign
def onJoinGroup(member):
group.state = PreparingRebalance
# signal EVERY member to drop EVERY partition
for m in group.members:
m.send(RevokeAll)
await m.onPartitionsRevoked_done
# group is idle here — nobody owns anything
plan = assignor.assign(group.members,
group.subscribed_topics)
group.generation += 1
for m, parts in plan.items():
m.send(SyncGroup(parts))
group.state = Stable
Consumer.onPartitionsRevoked
user callback — why every revoke costs wall-clock
def onPartitionsRevoked(partitions):
# everything below runs while the lane is BLACK
for p in partitions:
consumer.commitSync(offsets[p])
stateStore[p].flush()
stateStore[p].close()
# Streams: local RocksDB rebuilt on the next assignment
metrics.record('revoke.latency', now() - t0)
# only after this returns does Coordinator proceed

Where this sits in Build Kafka

Scene 07 of 13, in the Scale act — Rebalancing without halting every consumer.. Eager revokes everyone; cooperative-sticky only the lanes that move.

Up next. Exactly-once — when retries, restarts, and rebalances all stack up, what stops a duplicate or a zombie writer? Three independent monotonic counters do it.

All 13 scenes in Build Kafka · Every curriculum

Built with Arqly
Every scene in Build Kafka builds on the one before it.All 13 Build Kafka scenes