Skip to content
ConsensusKRaftDistributed Systems

Quorum Consensus: Decisions Without a Boss

How Beelogik commits state across a decentralized cluster without a central coordinator — and why quorum is the right answer for high-stakes workloads.

Beelogik Engineering 12 min read

Kafka gives us durable, ordered, replayable event streams. But before a record can be appended to a Kafka partition, the cluster has to agree on a smaller set of questions: who owns this partition, who is the leader, what topics exist, and where are the replicas.

Those are metadata decisions, and they require immediate, unambiguous, linearizable agreement. There is no “eventually” for “which broker is the leader of partition 7.” If two brokers each believe they lead the same partition, the log splits and records are lost.

This is the problem that consensus algorithms solve, and it is the problem we solve at Beelogik with quorum-based consensus, implemented in Kafka’s KRaft (Kafka Raft) controller.

What a quorum actually is

A quorum is a majority of the voting members of a cluster. In a five-node cluster, the quorum is three nodes. In a seven-node cluster, it is four. Any decision that requires a quorum must be acknowledged by at least that many nodes before it is considered committed.

The magic property of majority quorums is that any two quorums overlap.

Consider a five-node cluster: {A, B, C, D, E}. Any majority has at least three members. If you take two majorities, they must share at least one node — because 3 + 3 = 6, but there are only 5 nodes. By pigeonhole, at least one node is in both quorums.

This single arithmetic fact is the foundation of every safe consensus algorithm ever written. It is why consensus is possible at all in a system where nodes can fail and messages can be lost.

KRaft: Kafka’s consensus layer

For most of its history, Kafka stored its metadata in Apache ZooKeeper — an external service that provided a distributed filesystem and coordination primitives. ZooKeeper worked, but it was a second system to operate, a second set of failure modes to understand, and a second place where metadata could become inconsistent with itself.

KRaft — Kafka Raft — replaces ZooKeeper with a consensus protocol built into Kafka itself. A small set of Kafka nodes are designated as controllers; these controllers form a Raft group that manages all cluster metadata. Brokers no longer need to talk to ZooKeeper; they talk to the controller quorum.

The three problems KRaft solves are the same three problems any consensus algorithm solves.

Leader election

When a controller follower has not heard from the leader within an election timeout, it becomes a candidate and requests votes from the other controllers. A candidate wins if it receives votes from a majority of the controller set. Because every candidate in the same term requests votes from the same set, and because each node votes at most once per term, at most one candidate can win a given term.

If two candidates tie, no one wins, the term advances, and a new election runs. In practice, ties are rare because we randomize election timeouts — this is the key trick that makes Raft elections converge quickly.

Metadata replication

Once a controller leader is elected, it accepts metadata changes — new topics, partition reassignments, configuration updates — and replicates them to follower controllers. A change is committed when a majority of controllers have acknowledged it. At that point, the leader applies the change to its in-memory metadata state machine and returns success to the requesting broker.

The safety guarantee here is subtle but crucial: a committed metadata change is never lost, even if the controller leader crashes immediately after committing it. Why? Because the change is already present on a majority of controllers. When a new leader is elected, it must have been voted for by a majority — which must overlap with the majority that held the committed change — and Raft’s election restriction guarantees that a node will not vote for a candidate whose log is less complete than its own.

The overlap property again.

Log compaction

A growing metadata log is a growing problem: eventually it will not fit in memory, and replaying it on restart will take longer than the cluster can be down. KRaft solves this with snapshots — periodic checkpoints of the metadata state machine that allow the log before the snapshot to be discarded.

In Beelogik, snapshots are tiered: hot snapshots stay in memory, warm snapshots live on local NVMe, and cold snapshots are pushed to object storage. Recovery time is proportional to the tier the node recovers from, which lets us offer tiered RTO (recovery time objective) SLAs.

What quorum really costs

Every decision that goes through consensus pays a latency tax. In the worst case, a metadata change must travel from the requesting broker to the controller leader, from the leader to a majority of controller followers, and back. On a single-region cluster, this is sub-millisecond. Across regions, it is bounded by the speed of light.

We do not hide this tax. Instead, we let you choose where to pay it.

  • Regional consensus — the controller quorum lives entirely in one region. Latency is minimal, but a full region outage stops metadata changes (brokers continue serving produce and fetch requests using their existing assignment, so data-plane traffic is unaffected).
  • Cross-region consensus — the controller quorum is spread across regions. Latency is higher (bounded by the slowest link in the quorum), but the cluster survives the loss of an entire region.
  • Hybrid — critical metadata uses cross-region consensus; everything else uses regional consensus with asynchronous replication across regions.

Most of our customers start with regional consensus and migrate to the hybrid model as their compliance requirements tighten.

A worked example

Suppose a five-controller KRaft cluster is spread across three regions: Region A: controller-1 (leader), controller-2 (follower) Region B: controller-3 (follower), controller-4 (follower) Region C: controller-5 (follower)

A broker in Region A requests a partition reassignment. The controller leader replicates the change to all four followers. The quorum is three nodes total, including the leader, so the change commits as soon as any two followers acknowledge.

Region A fails. Regions B and C are still healthy — three controllers are still alive. A new controller leader is elected from Region B (a majority of the surviving three controllers is two, and Region B has two of them). The cluster continues accepting metadata changes.

For the change to be lost, both the leader and every controller that acknowledged it before the failure would have to die. That is the essence of the “no single point of failure” promise — not that nothing fails, but that no single failure is sufficient to lose committed state.

The performance envelope

In production, our KRaft controller path looks like this:

  • p50 metadata commit latency, regional: 0.8 ms
  • p99 metadata commit latency, regional: 3.2 ms
  • p50 metadata commit latency, cross-region (US-EU): 84 ms
  • p99 metadata commit latency, cross-region (US-EU): 118 ms

These numbers are for a five-controller cluster on NVMe storage. They are dominated by round-trip network time, not by the consensus algorithm itself. On a local network, Raft is essentially free — the algorithm is CPU-cheap and the log is append-only.

Note that these latencies apply to metadata operations (topic creation, partition reassignment). Data-plane operations (produce, fetch) go directly to the partition leader and are unaffected by controller latency.

Why this matters for your business

If you are running a system where “eventually” is not good enough for control-plane decisions — cluster membership, leader election, partition assignment — you need consensus. And if you need consensus at scale, you need quorum, because quorum is the only mechanism that gives you both safety and availability within the CAP theorem’s constraints.

What you get, at the product level, is:

  • Strong consistency for metadata. When a topic change returns success, it is committed — period.
  • Availability during partial failure. As long as a majority of controllers are alive, the cluster keeps accepting metadata changes.
  • No ZooKeeper to operate. KRaft removes an entire external dependency from the stack.
  • No coordinator. There is no central node whose loss would halt the system.

The beehive does not have a boss. Neither does a Beelogik KRaft group.