6.2 The controller
This is one of the most important distributed-systems ideas in Kafka.
2.1 Election via a single ephemeral node
The node is ephemeral, so “this way, we guarantee that the cluster will only have one controller at a time.”
On controller loss:
"When the controller broker is stopped or loses connectivity to ZooKeeper, the ephemeral node will disappear. This includes any scenario in which the ZooKeeper client used by the controller stops sending heartbeats to ZooKeeper for longer than
zookeeper.session.timeout.ms. When the ephemeral node disappears, other brokers... will be notified through the ZooKeeper watch... and will attempt to create the controller node themselves. The first node to create the new controller becomes the next controller, while the others receive 'node already exists' and re-create the watch on the new node."
2.2 ⚠️ The controller epoch — zombie fencing
This is one of the most important distributed-systems ideas in Kafka.
"Each time a controller is elected, it receives a new, higher controller epoch number through a ZooKeeper conditional increment operation. The brokers know the current controller epoch, and if they receive a message from a controller with an older number, they know to ignore it."
Why it's necessary — the zombie scenario:
At t4 broker A’s GC pause ends and it has no idea anything happened. It still believes it is the controller, and resumes sending commands stamped epoch = 7.
"The controller epoch in the message, which allows brokers to ignore messages from old controllers, is a form of zombie fencing." And: "The controller uses the epoch number to prevent a 'split brain' scenario where two nodes believe each is the current controller."
The generalizable lesson: you cannot prevent a paused process from waking up and acting stale. You can only make its stale actions rejectable. Monotonic epochs + conditional increment is how.
2.3 Controller startup cost
"When the controller first comes up, it has to read the latest replica state map from ZooKeeper before it can start managing the cluster metadata and performing leader elections. The loading process uses async APIs, and pipelines the read requests to ZooKeeper to hide latencies. But even so, in clusters with large numbers of partitions, the loading process can take SEVERAL SECONDS."
2.4 What the controller does on broker failure
- The new leader learns: start serving producer and consumer requests.
- The followers learn: start replicating from the new leader.
Broker startup is the mirror image, with one difference:
"the main difference is that all replicas in the broker start as FOLLOWERS and need to catch up to the leader before they are eligible to be elected as leaders themselves."
Summary from the book: "Kafka uses ZooKeeper's ephemeral node feature to elect a controller and to notify the controller when nodes join and leave the cluster. The controller is responsible for electing leaders among the partitions and replicas whenever it notices nodes join and leave. The controller uses the epoch number to prevent a 'split brain' scenario."
6.1 Cluster membership — ZooKeeper ephemeral nodes
Note the three causes lumped together: stopped, network partition, long GC pause.
6.3 KRaft — the Raft-based controller
That is the broker-level analogue of controller zombie fencing: a lagging broker can currently accept writes it has no right to accept, because it doesn't yet know it lost leaders…