Cluster Coordination: KRaft and ZooKeeper
How a Kafka cluster agrees on its own metadata: ZooKeeper's legacy role, the KRaft controller quorum that replaced it, and the migration between them. Interviewers ask because KRaft is now the default and every operator has to reason about controllers.
part ofApache Kafkaoverview, primer and where to startread it →on this pageshowhide
explore
- ZooKeeper's Historical Role5 questions
- ZooKeeper Scalability Limits5 questions
- KRaft Architecture and Metadata Quorum5 questions
- Raft Protocol and Leader Election5 questions
- Active Controller and Metadata Log5 questions
- Metadata Snapshots and Log Replay5 questions
- Broker Registration and Heartbeats5 questions
- Dynamic Quorum Reconfiguration5 questions
- ZooKeeper-to-KRaft Migration6 questions
- KRaft Storage Format and Cluster Bootstrap5 questions
- KRaft Operations and Observability5 questions
questions
61 · 12 sectionsBefore KRaft, what did Kafka use ZooKeeper for, and what kinds of cluster metadata lived in it?
basics
~20 sZooKeeper was an external coordination service that stored Kafka's cluster metadata: which brokers are alive, the list of topics and partitions, partition leaders and ISR, controller election, and ACLs/configs. Kafka depended on it to run.
Why was ZooKeeper deprecated and removed from Kafka (KIP-500)? What concrete problems did the ZK dependency cause?
basics
~20 sRunning ZooKeeper meant operating a second distributed system alongside Kafka — extra ops, a separate failure domain, and a metadata-scaling bottleneck (slow controller failover and startup at high partition counts). KIP-500 replaced it so Kafka manages its own metadata, simplifying operations and improving scalability.
What is an ephemeral znode, and how did Kafka use ephemeral znodes and ZooKeeper sessions to detect a dead broker?
basics
~20 sAn ephemeral znode is a ZooKeeper node that exists only while the client's session is alive; it auto-deletes when the session ends. Each broker created one under /brokers/ids. When the broker's session timed out, the znode vanished, signaling the broker was dead.
Walk through how the Kafka controller was elected via ZooKeeper, and what happened on controller failover.
basics
~20 sBrokers raced to create a single ephemeral znode at /controller; whoever created it first became the controller. If that broker died, its session expired, /controller was deleted, every other broker got a watch notification, and they raced again to elect a new controller.
What were the key ZooKeeper znode paths Kafka used (e.g. /brokers, /controller, /admin), and how could you inspect them?
basics
~10 sKafka kept its metadata under well-known ZK paths: /brokers/ids (live brokers), /brokers/topics (topic/partition assignment), /controller (current controller), /controller_epoch, /admin (admin operations like deletes/reassignments), /config, and /kafka-acl. You inspected them with zookeeper-shell.sh.
What was KIP-500 and what core problems with the ZooKeeper-based design did it set out to solve?
basics
~10 sKIP-500 is the Kafka proposal to remove the ZooKeeper dependency and manage metadata inside Kafka itself (KRaft). It targeted slow controller failover, scalability limits on partitions, operational complexity, and having two systems to run.
Why did Kafka's controller failover get slower as the number of partitions in the cluster grew, in the ZooKeeper-based architecture?
basics
~20 sOn the old design, when the controller failed the new one had to load ALL cluster metadata from ZooKeeper before it could act. More partitions meant more data to read, so recovery time grew with partition count.
What problem did having ZooKeeper AND the Kafka controller both hold metadata create, and how could they diverge?
basics
~20 sMetadata lived in ZooKeeper but the controller cached it in memory and also pushed it to brokers. With two copies that update separately, they could drift, causing brokers to act on stale or inconsistent metadata.
What is a ZooKeeper 'watch storm' in the context of Kafka, and why did it hurt at scale?
basics
~20 sKafka watched many ZooKeeper znodes to learn about changes. One event (like a broker dying) could fire huge numbers of watch notifications at once, flooding ZooKeeper and the controller with work and slowing metadata propagation.
As a principal engineer, explain why the ZooKeeper-based design effectively capped the number of partitions per cluster, tying together failover time, propagation latency, and change-rate limits.
basics
~20 sSeveral costs all grew with partition count: controller failover (O(partitions) ZK reload), metadata propagation to brokers, and the rate of changes ZooKeeper/watches could absorb. Together they made very large clusters slow to recover and risky, capping practical partition counts in the tens to low hundreds of thousands.
What is KRaft in Apache Kafka, and what problem was it introduced to solve?
basics
~10 sKRaft (Kafka Raft) is Kafka's built-in consensus protocol that replaces Apache ZooKeeper for storing cluster metadata. It lets Kafka manage its own metadata internally, so you no longer run a separate ZooKeeper cluster.
Describe the __cluster_metadata topic in KRaft: how it is structured, who reads and writes it, and why it is special.
basics
~20 s__cluster_metadata is the internal Kafka topic that holds all cluster metadata as an event log. It has a single partition replicated across the controller quorum via Raft. The active controller writes to it; controllers and brokers read it to stay in sync.
Explain the process.roles, node.id, and controller.quorum.voters configuration properties in KRaft. What does each control?
basics
~20 sprocess.roles sets whether a node is a broker, controller, or both (combined). node.id is the node's unique integer ID in the cluster. controller.quorum.voters lists the controller nodes (id@host:port) that form the metadata Raft quorum, so every node knows who the controllers are.
How does metadata propagate from the active controller to brokers in KRaft, and how does this differ from the ZooKeeper-era controller model?
basics
~20 sIn KRaft, brokers pull metadata changes by fetching the __cluster_metadata log from the active controller and applying records incrementally to a local metadata cache. In the ZooKeeper era, the controller pushed full LeaderAndIsr/UpdateMetadata RPCs to brokers, which was slower and harder to scale.
In a KRaft controller quorum, how is a metadata write committed, how many controller failures can the cluster tolerate, and what happens when quorum is lost?
basics
~20 sA metadata record is committed once a majority of controller voters have persisted it. With N voters the cluster tolerates floor((N-1)/2) failures (so 3 tolerate 1, 5 tolerate 2). If a majority is lost, no new metadata can be committed and the controller becomes read-only until quorum returns.
What is KRaft and what role does the Raft consensus protocol play in it?
basics
~20 sKRaft (Kafka Raft) is Kafka's built-in consensus mechanism that replaces ZooKeeper for storing cluster metadata. A small group of controllers uses a Raft-style protocol to elect a leader and agree on an ordered log of metadata changes.
How does leader election work in KRaft, and what is a leader epoch?
basics
~20 sWhen voters detect no leader, a candidate increases the epoch (a term counter), votes for itself, and sends Vote requests to peers. If a majority grants votes, it becomes leader for that epoch. The epoch is a monotonically increasing number that totally orders leadership periods.
How does KRaft decide a metadata record is committed, and how does the high watermark advance?
basics
~20 sA record is committed once it has been replicated (fetched) by a majority of the voters. The high watermark is the highest offset known to be on a majority; the leader advances it as Fetch offsets arrive, and only records at or below the HWM are visible/applied.
How does KRaft replicate the metadata log, and why is it described as pull-based rather than push-based Raft?
basics
~20 sIn KRaft, followers and observers pull records from the leader using Fetch requests, just like Kafka consumers pull from a partition leader. Classic Raft instead has the leader push records via AppendEntries. KRaft reuses Kafka's existing fetch/replication machinery.
How do you size a KRaft controller quorum, and what fault tolerance does each size give?
basics
~20 sUse an odd number of controllers, typically 3 or 5. A quorum of N tolerates floor((N-1)/2) failures: 3 tolerate 1, 5 tolerate 2. Odd counts give the best fault tolerance per node since a majority must still be reachable to elect a leader and commit.
In KRaft mode, what is the active controller and what role does it play in managing cluster metadata?
basics
~20 sThe active controller is the single elected controller node that is the only one allowed to write cluster metadata (topics, partitions, broker registrations) by appending records to a replicated log. Other nodes follow and replay that log.
Explain how brokers consume the metadata log as an event-sourced state machine, and what guarantees this provides.
basics
~20 sBrokers fetch the replicated metadata log and apply each record in order to build their in-memory view of the cluster. Because everyone replays the same ordered log, all nodes converge to the same state (eventual consistency), with the offset marking how far each has caught up.
Name some of the record types the active controller appends to the metadata log and explain what each represents.
basics
~10 sMetadata changes are encoded as typed records, e.g. TopicRecord (a topic was created), PartitionRecord/PartitionChangeRecord (a partition's replicas/leader/ISR), and RegisterBrokerRecord plus BrokerRegistrationChangeRecord (a broker joined or changed state).
Walk through what happens, in terms of metadata-log records, when a broker starts up and joins a KRaft cluster.
basics
~20 sThe broker sends a registration request to the active controller, which appends a RegisterBrokerRecord (broker id, endpoints, incarnation) to the log. The broker starts fenced, sends periodic heartbeats, and once caught up the controller appends an unfence record so it can host partition leaders.
What are the consistency and availability tradeoffs of routing all metadata writes through a single active controller and replicated log, compared to the old ZooKeeper model?
basics
~20 sThe single-writer log gives one ordered history of all metadata changes, so every node converges to the same state and failover is fast. The cost: writes need a controller-quorum majority, and the active controller is a serialization point for metadata throughput.
In a KRaft cluster, why does Kafka periodically take metadata snapshots of the __cluster_metadata log?
basics
~20 sThe metadata log grows forever as records are appended. Snapshots capture the current state at a point in time so old log records can be deleted, keeping the log small and making new nodes catch up faster.
Walk through how a KRaft broker or controller materializes the current cluster state at startup using snapshots and the log.
basics
~20 sIt loads the most recent snapshot to reach that offset instantly, then replays the metadata log records after the snapshot offset in order, applying each to its in-memory state, until it reaches the end of the log.
How do metadata snapshots and log replay enable fast controller failover in KRaft compared with the ZooKeeper-based architecture?
basics
~20 sStandby controllers continuously replay the same metadata log, so they already hold near-current state. When the active controller fails, a standby that is caught up can take over in seconds — no full metadata reload from ZooKeeper is needed.
What controls when KRaft generates a new metadata snapshot, and what does metadata.log.max.record.bytes.between.snapshots do?
basics
~10 sSnapshots are triggered by thresholds. metadata.log.max.record.bytes.between.snapshots sets how many bytes of new log records may accumulate since the last snapshot before a new one is generated. A time-based setting also forces periodic snapshots.
What happens when a KRaft node has fallen so far behind that the metadata log records it needs have already been truncated?
basics
~20 sThe leader cannot send records that were already deleted, so it transfers the latest snapshot to the lagging node instead. The node installs that snapshot to jump forward, then resumes fetching and replaying the log tail.
In a KRaft Kafka cluster, how does a broker tell the controller it is still alive, and what RPC is involved?
basics
~10 sEach broker periodically sends a BrokerHeartbeat request to the active controller. The interval is set by broker.heartbeat.interval.ms (default 2000 ms). If heartbeats stop, the controller eventually fences the broker.
Explain fencing vs unfencing of a broker in KRaft: what triggers each transition and what are the operational effects?
basics
~10 sA fenced broker is registered but excluded from leadership and ISRs. Brokers start fenced; they become unfenced once caught up to the metadata log and heartbeating normally. Missing heartbeats or controlled shutdown re-fences them.
What is a broker epoch in KRaft, how is it assigned, and what problem does it solve during broker restarts?
basics
~20 sA broker epoch is a monotonically increasing id the controller assigns at registration (the metadata log offset of the registration record). Every heartbeat carries it, so the controller can reject stale requests from a previous broker incarnation.
Walk through how controlled shutdown works for a broker in KRaft and why it matters for availability.
basics
~10 sOn graceful stop, the broker signals shutdown via its heartbeat. The controller moves leadership off it to in-sync replicas before it exits, avoiding abrupt leader elections and minimizing produce/consume disruption.
You operate a KRaft cluster in a cloud with occasional network jitter. How do you reason about tuning broker.heartbeat.interval.ms and broker.session.timeout.ms, and what are the failure modes at each extreme?
basics
~20 sKeep the session timeout a comfortable multiple of the heartbeat interval (defaults: 2000 ms / 9000 ms ~= 4.5x). Too short a timeout causes false fencing on jitter; too long delays detection of real failures.
Controller Quorum Sizing and Split-Brain Avoidance
all 5 Controller Quorum Sizing and Split-Brain Avoidance questions →How many controllers (voters) should a KRaft controller quorum have, and how many failures can each size tolerate?
basics
~10 sUse an odd number, usually 3 or 5 voters. A quorum needs a majority alive to work. 3 voters tolerate 1 failure; 5 voters tolerate 2 failures.
Why does adding a controller to make the quorum an even number (e.g. going from 3 to 4) add no extra fault tolerance?
basics
~20 sFault tolerance depends on the majority, and the majority jumps the same way. 3 voters need 2 and tolerate 1; 4 voters need 3 and still tolerate only 1 — but with one more node that can fail.
How do leader epochs and fencing prevent split-brain when a deposed KRaft controller (or a partition leader) comes back without realizing it was replaced?
basics
~20 sEvery leadership term has a number called the epoch that only increases. When a new leader is elected the epoch goes up. A stale old leader still on the lower epoch is rejected (fenced) by everyone, so it can't commit anything or be obeyed.
When the active KRaft controller fails, how does the quorum elect a new active controller and resume serving metadata?
basics
~20 sFollowers stop hearing heartbeats from the dead leader, a candidate starts a new election term, and whichever candidate gets votes from a majority of voters becomes the new active controller and continues the metadata log.
Should KRaft controllers run on the same nodes as brokers, and why do production deployments isolate them?
basics
~20 sIn production, run controllers as dedicated nodes (process.roles=controller), separate from brokers (process.roles=broker). Combined-mode nodes are fine for dev but mean broker load can starve the metadata quorum and a failure takes out both roles at once.
What is the difference between the static controller.quorum.voters config and the dynamic quorum introduced by KIP-853?
basics
~10 sStatic quorums list every controller in controller.quorum.voters on all nodes, fixed at startup. KIP-853 dynamic quorums (KRaft) let you add or remove controllers at runtime with AddRaftVoter/RemoveRaftVoter instead, using controller.quorum.bootstrap.servers.
Walk through how you add a new controller voter to a running dynamic KRaft cluster and how you remove one.
basics
~20 sFormat the new controller with --no-initial-controllers so it joins as an observer, start it, let it catch up on the metadata log, then run kafka-metadata-quorum add-controller (AddRaftVoter) to promote it to a voter. To remove one, run remove-controller (RemoveRaftVoter), then shut the node down.
What is directory.id in a dynamic KRaft quorum and why does each voter need one?
basics
~20 sdirectory.id is a UUID identifying a controller's metadata log directory. It uniquely tags each voter so the quorum can tell apart a node that was wiped and rejoined from the original, preventing a stale replica from being mistaken for a current voter.
In a dynamic KRaft cluster, what is the difference between an observer and a voter, and how does a node move between these roles?
basics
~20 sA voter is a controller in the quorum's voter set that votes in elections and counts toward majority commits. An observer replicates the metadata log but doesn't vote or count. A node joins as an observer, catches up, then add-controller promotes it to voter.
Explain joint consensus and how it makes single-voter membership changes safe in KRaft.
basics
~20 sJoint consensus is a transitional Raft state where decisions need a majority of BOTH the old and the new voter sets at once. This overlap guarantees no two leaders can be elected from disjoint majorities during a membership change, keeping the cluster safe.
What is the ZooKeeper-to-KRaft migration (KIP-866) and why does Kafka need it?
basics
~10 sIt is the online process that moves a Kafka cluster's metadata from ZooKeeper to KRaft (Kafka's built-in Raft controller) without downtime. Kafka needs it because ZooKeeper is deprecated and removed in Kafka 4.0.
How do you provision the KRaft controller quorum and enable migration mode for a KIP-866 migration?
basics
~10 sStart a new KRaft controller quorum with process.roles=controller and zookeeper.metadata.migration.enable=true, pointing it at the existing ZooKeeper. It must reuse the cluster's existing cluster.id and connect to the same ZK.
Explain the dual-write phase and how brokers are rolled from ZooKeeper mode to KRaft mode during migration.
basics
~20 sIn dual-write, the active KRaft controller writes every metadata change to both KRaft and ZooKeeper, keeping them in sync. Then brokers are restarted one at a time into KRaft mode (process.roles=broker, no zookeeper.connect) until all are migrated.
How do you finalize a KIP-866 migration and verify the cluster is fully on KRaft?
basics
~20 sAfter every broker is in KRaft mode, finalize by removing zookeeper.metadata.migration.enable (and zookeeper.connect) from the controllers and restarting them as a pure KRaft quorum. Verify via the migration/controller metrics showing the dual-write phase ended, then decommission ZooKeeper.
When can you roll back a KIP-866 migration, and what makes the rollback window close?
basics
~20 sYou can roll back any time during the dual-write phase, because ZooKeeper is still a current mirror. The window closes once you finalize the migration by taking the controllers out of migration mode; after that ZooKeeper is no longer written and you cannot return.
KRaft Storage Format and Cluster Bootstrap
all 5 KRaft Storage Format and Cluster Bootstrap questions →Why must you run kafka-storage.sh format before starting a KRaft-mode Kafka broker or controller, and what does it create?
basics
~20 sIn KRaft mode you must format each node's log directory first. kafka-storage.sh format writes a meta.properties file containing a shared cluster.id and the node.id, so the node knows which cluster it belongs to. Without it, the node refuses to start.
Show the kafka-storage.sh commands to bootstrap a KRaft cluster's storage and explain the key meta.properties fields and common formatting pitfalls.
basics
~10 sGenerate one cluster id with kafka-storage.sh random-uuid, then run kafka-storage.sh format -t <id> -c <config> on each node. That writes meta.properties holding cluster.id, node.id, and version. Use --ignore-formatted to skip already-formatted dirs.
Where does KRaft store the cluster metadata on disk, and what is the on-disk layout of the __cluster_metadata log directory?
basics
~10 sKRaft stores metadata in an internal single-partition topic, __cluster_metadata-0, written as a normal Kafka log (segments, indexes) plus snapshot files. By default it lives under log.dirs; you can move it with metadata.log.dir.
Walk through what happens during first-boot quorum formation in a KRaft cluster: how do controllers find each other, elect a leader, and reach the point where brokers can register?
basics
~20 sOn first boot each controller reads its formatted meta.properties (cluster.id, node.id) and the voter set (static voters config or bootstrap.checkpoint). Controllers contact each other on the controller listener, run a Raft election to pick an active controller (leader), and once a quorum agrees, brokers can register and the cluster is live.
Explain how --initial-controllers and the bootstrap.checkpoint mechanism work for first-boot quorum bootstrap in newer Kafka versions, and how that differs from controller.quorum.voters.
basics
~20 sOlder KRaft uses a static controller.quorum.voters list. Newer Kafka (KIP-853) supports dynamic quorums: at format time you pass --initial-controllers (or --standalone) so the formatter writes a bootstrap.checkpoint snapshot recording the initial voter set, and the quorum bootstraps from that instead of a static config line.
What does kafka-metadata-quorum.sh describe show, and how do you use it to check the health of a KRaft cluster?
basics
~10 skafka-metadata-quorum.sh describe reports the KRaft metadata quorum: who the current leader controller is, the list of voters and observers, and how far each replica lags behind the leader's metadata log.
What does the ActiveControllerCount metric mean in a KRaft cluster, and what value should it have?
basics
~10 sActiveControllerCount is a JMX gauge that is 1 on the controller node that is the current active (leader) controller and 0 on the others. Summed across the cluster it should always equal exactly 1.
Which controller metrics track how far the metadata log has been applied, and how do you use lastAppliedRecordOffset to detect a lagging or stuck controller?
basics
~20 sControllers expose lastAppliedRecordOffset (the highest metadata-log offset the node has applied to its in-memory state), plus lastAppliedRecordTimestamp and lastAppliedRecordLagMs. Compare a follower's applied offset/lag against the leader to spot a controller that is stuck or falling behind.
Design a monitoring and alerting strategy for KRaft quorum health. Which signals do you collect and what do you alert on?
basics
~20 sCollect controller JMX metrics (ActiveControllerCount, lastAppliedRecordOffset/LagMs, metadata error counts, election/epoch counters) plus periodic kafka-metadata-quorum.sh describe output. Alert on: cluster ActiveControllerCount sum != 1, any voter lagging or missing, rising applied lag, and non-zero metadata errors.
What is kafka-metadata-shell.sh used for, and when would you reach for it instead of kafka-metadata-quorum.sh?
basics
~10 skafka-metadata-shell.sh opens an interactive, filesystem-like shell over the contents of the KRaft __cluster_metadata log (or a snapshot). You browse the actual metadata records (topics, brokers, configs); kafka-metadata-quorum.sh instead reports quorum/replication health.