Learn Labs
4. Kafka Consumers: Reading Data from Kafka

4.2 Rebalancing — the mechanism, and why it hurts

Rebalance protocol

protocol
2
rebalance rounds
partitions during the rebalance
consumption stopped25%
Safe

Only 3 of 12 partitions pause (25%). The rest keep consuming throughout. The cost is two rebalance rounds instead of one, which is almost always the better trade.

A rebalance happens on every scale event and every deploy. Eager rebalancing revokes every partition from every consumer first, so the whole group stops; cooperative revokes only the partitions that actually change owner.

When does a rebalance happen?

  1. A new consumer joins the group → it starts consuming partitions previously owned by another consumer
  2. A consumer shuts down or crashes → it leaves; its partitions go to a remaining consumer
  3. The topics the group consumes are modified — e.g. an administrator adds new partitions

"Moving partition ownership from one consumer to another is called a rebalance. Rebalances are important because they provide the consumer group with high availability and scalability ... but in the normal course of events they can be fairly undesirable."

2.1 Eager rebalance — "stop the world"

revokerevokerevokeall rejoinC1p0 p1C2p2 p3C3p4 p5C1--C2--C3--C1p0 p3C2p1 p4C3p2 p5PHASE 1: everyone gives up EVERYTHINGALL CONSUMPTION STOPPED —entire group unavailablePHASE 2: brand-newassignmentconsumption resumes

“During an eager rebalance, all consumers stop consuming, give up their ownership of all partitions, rejoin the consumer group, and get a brand-new partition assignment. This is essentially a short window of unavailability of the entire consumer group.”

Figure 4.2.12.1 Eager rebalance — 'stop the world'

"During an eager rebalance, all consumers stop consuming, give up their ownership of all partitions, rejoin the consumer group, and get a brand-new partition assignment. This is essentially a short window of unavailability of the entire consumer group. The length of the window depends on the size of the consumer group as well as on several configuration parameters."

2.2 Cooperative (incremental) rebalance

give up p1give up p3keeps bothC1p0 p1C2p2 p3C3p4 p5C1p0C2p2C3p4 p5C4p1 p3PHASE 1: leader tells consumers which SUBSET they'll lose(new consumer C4 is joining)p0, p2, p4, p5 KEEP FLOWINGPHASE 2: the leader assignsthe now-orphaned partitions
  • Only the REASSIGNED partitions paused. No “stop the world.”
  • “May take a few iterations until a stable assignment is achieved.”
Figure 4.2.22.2 Cooperative (incremental) rebalance

"Cooperative rebalances (also called incremental rebalances) typically involve reassigning only a small subset of the partitions from one consumer to another, and allowing consumers to continue processing records from all the partitions that are not reassigned. ... This is especially important in large consumer groups where rebalances can take a significant amount of time."

Sequence: (1) group leader informs all consumers they will lose ownership of a subset; consumers stop consuming those and give up ownership. (2) Group leader assigns the now-orphaned partitions to new owners.

2.3 Heartbeats, the group coordinator, and the two death-detection paths

heartbeatsGROUP COORDINATORa designated broker — can differ per groupConsumer C1Consumer C2Consumer C3BACKGROUND THREADsends heartbeatsMAIN THREADpoll() loopC2 and C3 hold the same two threads.

“As long as the consumer is sending heartbeats at regular intervals, it is assumed to be alive. If the consumer stops sending heartbeats for long enough, its session will timeout and the group coordinator will consider it dead and trigger a rebalance.”

Figure 4.2.32.3 Heartbeats, the group coordinator, and the two death-detection paths

"As long as the consumer is sending heartbeats at regular intervals, it is assumed to be alive. If the consumer stops sending heartbeats for long enough, its session will timeout and the group coordinator will consider it dead and trigger a rebalance."

The two failure paths — and why there must be two:

PathDetectsConfigWhy it exists
Session timeout (missed heartbeats)The process is dead / unreachablesession.timeout.ms + heartbeat.interval.msThe obvious case
Poll interval timeoutThe main thread is stuck while the background heartbeat thread is still happily heartbeatingmax.poll.interval.ms"There is a possibility that the main thread consuming from Kafka is deadlocked, but the background thread is still sending heartbeats. This means that records from partitions owned by this consumer are not being processed."

Crash vs clean shutdown:

CRASHno heartbeatscoordinator waitssession.timeout.msNO MESSAGES PROCESSEDfrom the dead consumer's partitionsclean close()consumer NOTIFIES the coordinatorrebalance IMMEDIATELYno session timeout to wait outthe gap shrinks“reducing the gap in processing”

The coordinator needs “a few seconds without heartbeats to decide it is dead,” and during those seconds no messages are processed from the dead consumer's partitions.

Figure 4.2.4Crash vs clean shutdown

This is why close() matters operationally, not just for tidiness. See §7.

2.4 How partition assignment actually works

  1. Consumer sends JoinGroup request to the group coordinator
  2. The first consumer to join becomes the group leader
  3. Leader receives, from the coordinator, the list of All consumers in the group (all that sent a heartbeat recently = considered alive)
  4. Leader runs an implementation of PartitionAssignor to decide which partitions go to which consumer
  5. Leader sends the assignment list to the GroupCoordinator
  6. Coordinator sends the info to all consumers
  7. Each consumer only sees its own assignment — the leader is the only client process with the full picture
  8. This whole process repeats on Every rebalance

Design note worth appreciating: assignment logic runs on a client, not the broker. That's why you can plug in your own PartitionAssignor without touching the cluster — and why the leader's Kafka client version determines which strategies are available.


On this page