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…
Timeline: work started 2019; preview in Kafka 2.8; "The Apache Kafka 3.0 release, planned for mid 2021, will include the first production version of KRaft, and Kafka clusters will be able to run with either the traditional ZooKeeper-based controller or KRaft."
3.1 The four motivating problems — each is a real production pain
"Kafka's existing controller already underwent several rewrites, but... it became clear that the existing model will not scale to the number of partitions we want Kafka to support."
- Inconsistent metadata, hard to detect. “Metadata updates are written to ZooKeeper synchronously but are sent to brokers asynchronously. In addition, receiving updates from ZooKeeper is asynchronous. All this leads to edge cases where metadata is inconsistent between brokers, controller, and ZooKeeper. These cases are challenging to detect.”
- Controller restart is a scaling bottleneck. “Whenever the controller is restarted, it has to read all the metadata for all brokers and partitions from ZooKeeper and then send this metadata to all brokers. Despite years of effort, this remains a major bottleneck — as the number of partitions and brokers increases, restarting the controller becomes slower.”
- Confused metadata ownership. “some operations were done via the controller, others via any broker, and others directly on ZooKeeper.” — which is exactly why Ch. 5 says never touch ZooKeeper directly.
- Two distributed systems to operate. “ZooKeeper is its own distributed system, and, just like Kafka, it requires some expertise to operate. Developers who want to use Kafka therefore need to learn two distributed systems, not just one.”
3.2 What has to be replaced
ZooKeeper today does two things:
- elect a controller
- store cluster metadata — registered brokers, configuration, topics, partitions, replicas
The controller itself manages metadata:
- elect leaders
- create and delete topics
- reassign replicas
⇒ All of this must be replaced.
3.3 The core idea — eat your own dog food
"The core idea behind the new controller design is that Kafka itself has a log-based architecture, where users represent state as a stream of events. The benefits are well understood: multiple consumers can quickly catch up to the latest state by replaying events. The log establishes a clear ordering between events and ensures that the consumers always move along a single timeline. The new controller architecture brings the same benefits to the management of Kafka's metadata."
- Brokers track the OFFSET of the latest metadata change fetched, and request only newer updates.
- Brokers persist metadata to disk → “start up quickly, even with millions of partitions.”
Key changes, one by one:
| Old | New |
|---|---|
| ZooKeeper elects the controller | "Using the Raft algorithm, the controller nodes will elect a leader from among themselves, without relying on any external system." |
| Controller failover = lengthy full metadata reload | "Because the controllers will now all track the latest state, controller failover will NOT require a lengthy reloading period in which we transfer all the state to the new controller." |
| Controller pushes updates to brokers | "brokers will FETCH updates from the active controller via a new MetadataFetch API" — offset-tracked, incremental, just like a normal fetch |
| Brokers rebuild metadata at startup | "Brokers will persist the metadata to disk" → fast startup at millions of partitions |
| Broker registration is ephemeral | "Brokers will register with the controller quorum and will remain registered until unregistered by an admin, so once a broker shuts down, it is offline but still registered." |
| A stale broker can serve requests | NEW FENCED STATE — see below |
3.4 The new fenced state — closing a real correctness hole
"Brokers that are online but are not up-to-date with the latest metadata will be FENCED and will not be able to serve client requests. The new fenced state will prevent cases where a client produces events to a broker that is no longer a leader but is too out-of-date to be aware that it isn't a leader."
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 leadership.
3.5 Migration strategy
"As part of the migration to the controller quorum, all operations that previously involved either clients or brokers communicating directly to ZooKeeper will be routed via the controller. This will allow seamless migration by replacing the controller without having to change anything on any broker."
Reference KIPs:
| KIP | Content |
|---|---|
| KIP-500 | Overall design of the new architecture |
| KIP-595 | How the Raft protocol was adapted for Kafka |
| KIP-631 | Controller quorum design, controller configuration, new CLI for cluster metadata |