skip to content

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 pageshow

explore

questions

61 · 12 sections

Before KRaft, what did Kafka use ZooKeeper for, and what kinds of cluster metadata lived in it?

level: juniorimportance: must knowfreq 70%
basics
~20 s

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

open as a page

Why was ZooKeeper deprecated and removed from Kafka (KIP-500)? What concrete problems did the ZK dependency cause?

level: seniorimportance: must knowfreq 75%
basics
~20 s

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

open as a page

What is an ephemeral znode, and how did Kafka use ephemeral znodes and ZooKeeper sessions to detect a dead broker?

level: middleimportance: should knowfreq 55%
basics
~20 s

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

open as a page

Walk through how the Kafka controller was elected via ZooKeeper, and what happened on controller failover.

level: seniorimportance: should knowfreq 50%
basics
~20 s

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

open as a page

What were the key ZooKeeper znode paths Kafka used (e.g. /brokers, /controller, /admin), and how could you inspect them?

level: middleimportance: nice to knowfreq 35%
basics
~10 s

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

open as a page

What was KIP-500 and what core problems with the ZooKeeper-based design did it set out to solve?

level: juniorimportance: must knowfreq 70%
basics
~10 s

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

open as a page

Why did Kafka's controller failover get slower as the number of partitions in the cluster grew, in the ZooKeeper-based architecture?

level: middleimportance: must knowfreq 62%
basics
~20 s

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

open as a page

What problem did having ZooKeeper AND the Kafka controller both hold metadata create, and how could they diverge?

level: seniorimportance: should knowfreq 45%
basics
~20 s

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

open as a page

What is a ZooKeeper 'watch storm' in the context of Kafka, and why did it hurt at scale?

level: seniorimportance: should knowfreq 38%
basics
~20 s

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

open as a page

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.

level: principalimportance: should knowfreq 30%
basics
~20 s

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

open as a page

What is KRaft in Apache Kafka, and what problem was it introduced to solve?

level: juniorimportance: must knowfreq 80%
basics
~10 s

KRaft (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.

open as a page

Describe the __cluster_metadata topic in KRaft: how it is structured, who reads and writes it, and why it is special.

level: middleimportance: must knowfreq 65%
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.

open as a page

Explain the process.roles, node.id, and controller.quorum.voters configuration properties in KRaft. What does each control?

level: middleimportance: must knowfreq 70%
basics
~20 s

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

open as a page

How does metadata propagate from the active controller to brokers in KRaft, and how does this differ from the ZooKeeper-era controller model?

level: seniorimportance: should knowfreq 45%
basics
~20 s

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

open as a page

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?

level: seniorimportance: should knowfreq 55%
basics
~20 s

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

open as a page

What is KRaft and what role does the Raft consensus protocol play in it?

level: juniorimportance: must knowfreq 70%
basics
~20 s

KRaft (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.

open as a page

How does leader election work in KRaft, and what is a leader epoch?

level: middleimportance: must knowfreq 65%
basics
~20 s

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

open as a page

How does KRaft decide a metadata record is committed, and how does the high watermark advance?

level: seniorimportance: must knowfreq 50%
basics
~20 s

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

open as a page

How does KRaft replicate the metadata log, and why is it described as pull-based rather than push-based Raft?

level: seniorimportance: must knowfreq 55%
basics
~20 s

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

open as a page

How do you size a KRaft controller quorum, and what fault tolerance does each size give?

level: principalimportance: should knowfreq 45%
basics
~20 s

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

open as a page

In KRaft mode, what is the active controller and what role does it play in managing cluster metadata?

level: juniorimportance: must knowfreq 70%
basics
~20 s

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

open as a page

Explain how brokers consume the metadata log as an event-sourced state machine, and what guarantees this provides.

level: seniorimportance: must knowfreq 50%
basics
~20 s

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

open as a page

Name some of the record types the active controller appends to the metadata log and explain what each represents.

level: middleimportance: should knowfreq 45%
basics
~10 s

Metadata 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).

open as a page

Walk through what happens, in terms of metadata-log records, when a broker starts up and joins a KRaft cluster.

level: seniorimportance: should knowfreq 35%
basics
~20 s

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

open as a page

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?

level: principalimportance: should knowfreq 30%
basics
~20 s

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

open as a page

In a KRaft cluster, why does Kafka periodically take metadata snapshots of the __cluster_metadata log?

level: juniorimportance: must knowfreq 62%
basics
~20 s

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

open as a page

Walk through how a KRaft broker or controller materializes the current cluster state at startup using snapshots and the log.

level: middleimportance: must knowfreq 55%
basics
~20 s

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

open as a page

How do metadata snapshots and log replay enable fast controller failover in KRaft compared with the ZooKeeper-based architecture?

level: seniorimportance: should knowfreq 38%
basics
~20 s

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

open as a page

What controls when KRaft generates a new metadata snapshot, and what does metadata.log.max.record.bytes.between.snapshots do?

level: seniorimportance: should knowfreq 40%
basics
~10 s

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

open as a page

What happens when a KRaft node has fallen so far behind that the metadata log records it needs have already been truncated?

level: principalimportance: nice to knowfreq 24%
basics
~20 s

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

open as a page

In a KRaft Kafka cluster, how does a broker tell the controller it is still alive, and what RPC is involved?

level: juniorimportance: must knowfreq 55%
basics
~10 s

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

open as a page

Explain fencing vs unfencing of a broker in KRaft: what triggers each transition and what are the operational effects?

level: middleimportance: must knowfreq 50%
basics
~10 s

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

open as a page

What is a broker epoch in KRaft, how is it assigned, and what problem does it solve during broker restarts?

level: seniorimportance: should knowfreq 35%
basics
~20 s

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

open as a page

Walk through how controlled shutdown works for a broker in KRaft and why it matters for availability.

level: seniorimportance: should knowfreq 30%
basics
~10 s

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

open as a page

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?

level: principalimportance: nice to knowfreq 18%
basics
~20 s

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

open as a page

How many controllers (voters) should a KRaft controller quorum have, and how many failures can each size tolerate?

level: juniorimportance: must knowfreq 70%
basics
~10 s

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

open as a page

Why does adding a controller to make the quorum an even number (e.g. going from 3 to 4) add no extra fault tolerance?

level: middleimportance: must knowfreq 55%
basics
~20 s

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

open as a page

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?

level: seniorimportance: must knowfreq 50%
basics
~20 s

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

open as a page

When the active KRaft controller fails, how does the quorum elect a new active controller and resume serving metadata?

level: seniorimportance: should knowfreq 45%
basics
~20 s

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

open as a page

Should KRaft controllers run on the same nodes as brokers, and why do production deployments isolate them?

level: seniorimportance: should knowfreq 40%
basics
~20 s

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

open as a page

What is the difference between the static controller.quorum.voters config and the dynamic quorum introduced by KIP-853?

level: juniorimportance: must knowfreq 60%
basics
~10 s

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

open as a page

Walk through how you add a new controller voter to a running dynamic KRaft cluster and how you remove one.

level: seniorimportance: must knowfreq 40%
basics
~20 s

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

open as a page

What is directory.id in a dynamic KRaft quorum and why does each voter need one?

level: middleimportance: should knowfreq 35%
basics
~20 s

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

open as a page

In a dynamic KRaft cluster, what is the difference between an observer and a voter, and how does a node move between these roles?

level: middleimportance: should knowfreq 25%
basics
~20 s

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

open as a page

Explain joint consensus and how it makes single-voter membership changes safe in KRaft.

level: principalimportance: should knowfreq 20%
basics
~20 s

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

open as a page

What is the ZooKeeper-to-KRaft migration (KIP-866) and why does Kafka need it?

level: juniorimportance: must knowfreq 70%
basics
~10 s

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

open as a page

How do you provision the KRaft controller quorum and enable migration mode for a KIP-866 migration?

level: middleimportance: must knowfreq 55%
basics
~10 s

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

open as a page

Explain the dual-write phase and how brokers are rolled from ZooKeeper mode to KRaft mode during migration.

level: seniorimportance: must knowfreq 50%
basics
~20 s

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

open as a page

How do you finalize a KIP-866 migration and verify the cluster is fully on KRaft?

level: middleimportance: should knowfreq 35%
basics
~20 s

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

open as a page

When can you roll back a KIP-866 migration, and what makes the rollback window close?

level: seniorimportance: should knowfreq 40%
basics
~20 s

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

open as a page

Why must you run kafka-storage.sh format before starting a KRaft-mode Kafka broker or controller, and what does it create?

level: juniorimportance: must knowfreq 70%
basics
~20 s

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

open as a page

Show the kafka-storage.sh commands to bootstrap a KRaft cluster's storage and explain the key meta.properties fields and common formatting pitfalls.

level: middleimportance: must knowfreq 55%
basics
~10 s

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

open as a page

Where does KRaft store the cluster metadata on disk, and what is the on-disk layout of the __cluster_metadata log directory?

level: middleimportance: should knowfreq 45%
basics
~10 s

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

open as a page

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?

level: seniorimportance: should knowfreq 28%
basics
~20 s

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

open as a page

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.

level: seniorimportance: should knowfreq 30%
basics
~20 s

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

open as a page

What does kafka-metadata-quorum.sh describe show, and how do you use it to check the health of a KRaft cluster?

level: juniorimportance: must knowfreq 70%
basics
~10 s

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

open as a page

What does the ActiveControllerCount metric mean in a KRaft cluster, and what value should it have?

level: middleimportance: must knowfreq 60%
basics
~10 s

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

open as a page

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?

level: seniorimportance: should knowfreq 40%
basics
~20 s

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

open as a page

Design a monitoring and alerting strategy for KRaft quorum health. Which signals do you collect and what do you alert on?

level: principalimportance: should knowfreq 30%
basics
~20 s

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

open as a page

What is kafka-metadata-shell.sh used for, and when would you reach for it instead of kafka-metadata-quorum.sh?

level: middleimportance: nice to knowfreq 25%
basics
~10 s

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

open as a page