skip to content

questions

6

In a Raft or Paxos cluster of N nodes, a write is considered committed once it has been acknowledged by a quorum. Why does that quorum have to be a strict majority (more than N/2) rather than some fixed count like 2 nodes, regardless of cluster size?

level: juniorimportance: must knowfreq 70%

answer

  1. majority quorums always intersect
  2. ⌊N/2⌋+1 formula
  3. pigeonhole: 3+3>5
  4. odd cluster sizes (3,5,7)
  5. minority partition halts, doesn't diverge

basics

~10 s

A majority is needed so any two quorums always share at least one node - that overlap is what stops two different decisions from being made at once.

solid answer

~40 s

Consensus needs any two quorums to intersect, so majority quorums (⌊N/2⌋+1) guarantee overlap regardless of which nodes are chosen. That overlap means if one majority commits value A, no other majority can commit value B without including a node that already saw A and refuses to accept a conflicting value - that's what prevents split-brain writes. A fixed count like 2 fails as clusters grow: in a 7-node cluster, two disjoint groups of 2 could each proceed independently and diverge. Majority quorums also tolerate ⌊(N-1)/2⌋ failures while staying available, giving the best fault tolerance for a given N. The trade-off is cost: larger clusters need more replicas to ack every write, raising latency and network chatter, which is why odd-sized clusters (3, 5, 7) are typical - they maximize fault tolerance per replica.

go deeper

for a junior

Should know a quorum means 'more than half must agree' and that this is why odd numbers of nodes (3, 5) are common; doesn't need to derive the intersection math.

for a middle

Should be able to state the quorum formula and explain, with a small example (e.g. 5 nodes, quorum of 3), why two majority quorums must overlap.

for a senior

Should connect quorum intersection directly to safety guarantees, explain the availability cost when a majority is unreachable, and reason about even- vs odd-sized clusters.

for a principal

Should discuss quorum reconfiguration hazards during membership changes (joint consensus), non-voting replica designs, and the latency/fault-tolerance trade-off curve as cluster size grows.

## What a quorum is A **quorum** is the minimum number of nodes in a cluster that must agree before an operation - typically a write or a leader election - is considered decided. In majority-based consensus protocols like Raft, Multi-Paxos, and ZooKeeper's ZAB, that quorum size is set to strict majority: for a cluster of N nodes, a quorum is `⌊N/2⌋ + 1`. The reason is a property called **quorum intersection**: - Any two majority subsets of the same set of N nodes must share at least one common member. - If N=5, a majority is 3; pick any two groups of 3 out of 5, and by the pigeonhole principle they cannot be disjoint, because 3+3=6 > 5. - That single overlapping node is the mechanism that prevents two conflicting decisions from both being finalized. ## How Raft and Paxos use the overlap Concretely, in Raft, a leader appends a new log entry locally then sends `AppendEntries` RPCs to all followers in parallel. It waits until a majority of the cluster (including itself) has persisted that entry, then marks it committed and applies it to the state machine. Because every future quorum must overlap with every past quorum, any node that later tries to get elected leader or accept a conflicting entry at the same log index will run into at least one node that already holds the committed entry - correctly-implemented followers refuse to overwrite committed history, and election rules (vote only for a candidate whose log is at least as up to date as mine) stop a lagging node from winning leadership in the first place. Paxos achieves the same guarantee more explicitly through **ballot numbers**: an acceptor promises not to accept any proposal with a lower ballot number than one it has already promised, so a proposer must first learn about any previously accepted value from a majority before it can safely propose - the intersection with the earlier accepting quorum guarantees it will see that value. ## Why the design exists The reason this design exists is that a distributed system has no single reliable source of truth - nodes can crash, messages can be delayed or lost, and the network can partition. Consensus protocols need a way to make a decision durable using only local disks and message-passing, without ever assuming a node 'knows' the global state. Quorums convert that problem into simple arithmetic: as long as you always require overlap between the readers/writers of a value, you get agreement without a central coordinator holding the only copy. ## The trade-off The trade-off is **availability versus consistency cost**. - A quorum protocol keeps working as long as a majority of nodes are reachable, but it stops making progress the moment a majority is unreachable - even if a large minority is up and healthy. - A 5-node cluster tolerates 2 simultaneous failures but not 3, even though 3 nodes could theoretically still talk to each other and to clients; the protocol deliberately refuses to let that minority make decisions, because it cannot be sure another minority isn't doing the same thing on the other side of a partition. That's the essence of choosing consistency (safety) over availability (liveness) under partition. - There's also a raw performance cost: every write must round-trip to a majority before it's acknowledged, so latency scales with the slower nodes in the quorum, and larger clusters (7 or 9 nodes) trade better fault tolerance for higher write latency and more network traffic, which is why production clusters are usually kept small (3 or 5) and read/scale traffic is handled by non-voting replicas instead of growing the voting set. ## Failure modes The most common failure mode from getting quorum sizing wrong is an even-sized cluster: 1. **Even sizes.** A 4-node cluster has a majority of 3, tolerating only 1 failure - the same fault tolerance as a 3-node cluster but with one more node to pay for and coordinate, which is why odd sizes are the standard recommendation. 2. **Quorum reconfiguration during membership changes** is a subtler failure mode: naively swapping the old member list for a new one can create a window where an old-majority quorum and a new-majority quorum don't overlap, producing split-brain; Raft's joint-consensus (or single-server-change) approach and Paxos's reconfiguration protocols exist specifically to keep old and new quorums intersecting throughout the transition. ## Where it shows up In production: - **etcd** (Raft-based) uses this exact majority-quorum mechanism to back Kubernetes' cluster state - a 3-node etcd cluster tolerates one node failure and requires 2 acks per write. - **Google's Chubby** lock service and **Apache ZooKeeper** (ZAB, a Paxos-like protocol) apply the same majority-quorum arithmetic to coordinate distributed locks and configuration. In every case, the majority requirement is what makes 'the cluster agreed on X' a statement you can trust even though no individual node can see the whole system at once.

  • What happens to a 5-node Raft cluster if a network partition splits it into a group of 3 and a group of 2?
    The group of 3 still has a majority, so it can elect a leader and keep committing writes normally. The group of 2 cannot reach a majority, so it can't elect a leader or commit anything - it simply stalls until the partition heals, which is what prevents the two sides from diverging.
  • Why do production clusters typically use 3 or 5 nodes rather than 4 or 6?
    Even-sized clusters pay for an extra node without gaining fault tolerance - a 4-node cluster still only tolerates 1 failure, same as 3 nodes, because its majority is 3. Odd sizes maximize the failures tolerated per node paid for, and adding nodes beyond 5 or 7 mostly adds write latency rather than meaningfully more availability.
  • Does adding a 6th non-voting read replica to a 5-node Raft cluster change the quorum size?
    No - non-voting members don't count toward quorum arithmetic; the majority is still computed over the 5 voting members. Read replicas are added specifically to scale read throughput without paying the write-latency and quorum-recalculation cost of adding voting members.

Like a company bylaw requiring board decisions to have signatures from more than half the board members: any two such signed decisions must share at least one signer, so that signer can catch and block a contradictory decision - a trick a fixed 'any 2 signatures' rule wouldn't guarantee once the board grows past 4.

saying these in an interview costs you the question

  • says quorum just means 'more than half the nodes are online', missing that it's about overlap between decision-making groups
  • claims a fixed quorum of 2 works for any cluster size
  • thinks adding more nodes always increases availability
  • doesn't know a minority partition should stop making progress, not keep serving writes
  • confuses quorum arithmetic with leader-election mechanics

context

open as a page

Walk through how a Raft leader replicates a new client write to its followers and decides when that entry is safely 'committed.'

level: middleimportance: must knowfreq 75%

basics

~10 s

The leader adds the write to its own list first, sends copies to the other servers, and once more than half of them have saved it, the leader tells everyone it's official.

open as a page

In the context of consensus protocols like Raft, what's the difference between a 'safety' property and a 'liveness' property, and give an example of each being satisfied even when the other is temporarily violated?

level: middleimportance: must knowfreq 55%

basics

~20 s

Safety means 'nothing bad ever happens' (like two leaders both thinking they won), and liveness means 'something good eventually happens' (like a write eventually succeeding). A system can stay safe while briefly failing to make progress.

open as a page

Suppose a network partition isolates the current Raft leader from the rest of the cluster, but that leader doesn't crash - it just can't reach anyone. The majority side elects a new leader. When the partition heals, how does the protocol prevent both leaders from being treated as valid, and what stops the old leader from having accepted writes in the meantime?

level: seniorimportance: must knowfreq 60%

basics

~20 s

Every election bumps a counter called a term. The old leader's term is now lower than the new one's, so once they reconnect, everyone - including the old leader - recognizes the higher term and steps down; and the old leader couldn't commit any writes alone anyway since it needs majority acks it never got.

open as a page

A distributed system needs multiple services to agree on whether a single business transaction succeeds or fails, atomically. When would you reach for a two-phase-commit-style atomic commit protocol instead of running the operation through a full consensus protocol like Raft or Paxos, and what do you give up either way?

level: seniorimportance: should knowfreq 45%

basics

~20 s

2PC is for getting several different services to agree on one all-or-nothing transaction right now; consensus (Raft/Paxos) is for keeping several copies of the same data in agreement over time, and it survives a coordinator crashing, while plain 2PC can freeze if its coordinator dies mid-transaction.

open as a page

Raft is often described as 'Paxos, but designed to be understandable.' Concretely, what problem does Raft's strong-leader approach solve relative to classic Multi-Paxos, and what does that design choice cost?

level: principalimportance: nice to knowfreq 25%

basics

~20 s

Multi-Paxos lets any node propose changes, which gets confusing and can cause competing proposals; Raft forces all changes through one leader at a time, which is much easier to reason about and implement correctly, but makes the leader a bottleneck and needs a full election protocol to replace it.

open as a page