Learn Labs
6. Kafka Internals

6.2 The controller

This is one of the most important distributed-systems ideas in Kafka.

2.1 Election via a single ephemeral node

broker AZooKeeperbroker Bbroker Ccreate /controllercreated ► CONTROLLERcreate /controller“node already exists”watch /controllercreate /controller“node already exists”watch /controller

The node is ephemeral, so “this way, we guarantee that the cluster will only have one controller at a time.”

Figure 6.2.12.1 Election via a single ephemeral node

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:

t0 · controller = broker Aepoch = 7t1 · broker A enters a long GC pauset2 · A’s ZK session times out/controller ephemeral node vanishest3 · broker B wins the racecontroller = broker B, epoch = 8t4 · A resumes, stamping epoch = 7broker A is now a ZOMBIEt5 · other brokers compare 7 to 8and DISCARD the messages

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.

Figure 6.2.2Why it's necessary — the zombie scenario

"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

controller notices a broker leftthe ZooKeeper watch, or a ControlledShutdownRequestevery partition led there needs a new leaderdetermine the new leader“simply the NEXT REPLICA IN THE REPLICA LIST of that partition”persist the new state to ZooKeeperpipelined async requests, to reduce latencyLeaderAndISR → brokers holding replicasbatched — one request carries many partitions on the same brokerUpdateMetadata → ALL brokersevery broker has a MetadataCache that must be refreshed
  • The new leader learns: start serving producer and consumer requests.
  • The followers learn: start replicating from the new leader.
Figure 6.2.32.4 What the controller does on broker failure

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."


On this page