Replicated systems keep copies of data on several machines so that some can fail. The moment you do that, you need a rule for how many copies an operation must touch before it counts. Touch all of them and one failure stops you. Touch one and two clients can read and write disjoint copies and never see each other. A quorum system is the rule in between: a collection of subsets of nodes, called quorums, chosen so that any two operations that must see each other's effects are guaranteed to meet on at least one node.
Most engineers meet quorums as the Dynamo-style formula for tuning reads and writes, which the R, W and N guide covers along with sloppy quorums. This article steps back and treats quorum systems as a design space. It starts from the intersection property, then builds majorities, weighted votes, grids, the asymmetric quorums of Flexible Paxos and the larger quorums Byzantine faults require, with exact availability numbers and a small Python checker you can use on your own configurations.
The one property: intersection
Suppose a write completes after it has been stored on the nodes of quorum A, and a later read consults the nodes of quorum B. If A and B share a node, the read sees the write, or at least sees that something newer exists and can go and find it. If they do not share a node, the read can return stale data and nothing in the protocol will notice. So the defining property of a quorum system is: every quorum that must observe an operation intersects every quorum that performs it.
Everything else is a choice about which pairs must intersect and how large the quorums are. Two measures drive that choice. Availability is the probability that some quorum is fully alive, given each node is up with probability p. Load is how often the busiest node is involved in an operation when quorums are chosen as evenly as possible; low load means throughput scales as nodes are added. Naor and Wool showed in 1998 that these pull against each other: quorum systems with very low load cannot also have the best availability.
Majority quorums
The simplest system with the intersection property takes every set of more than half the nodes. Two majorities of n nodes each hold more than n/2 nodes, so together they hold more than n, and must share one. With n = 2f + 1 nodes, a majority is f + 1, and the system survives any f failures.
Majorities are the availability champion: among quorum systems in which every pair intersects, a majority over an odd number of equally reliable nodes is the most available, as long as each node is up more than half the time. Raft and Multi-Paxos use them for both electing a leader and committing log entries; see Raft log replication. Their weakness is load. Every operation touches more than half the cluster, so adding nodes adds fault tolerance but not throughput. Even counts are a trap: four nodes need three for a majority and tolerate one failure, exactly like three nodes, while adding one more machine that can fail.
Weighted voting
Gifford's 1979 weighted voting generalises majorities. Each replica i gets vi votes, the total is V, a read needs r votes and a write needs w votes, with r + w > V so reads meet writes and 2w > V so two writes meet each other and can be ordered. Weights let you express that nodes are not equal. A replica in the primary data centre can carry more votes than one across an ocean; a witness node that stores only metadata can carry a single vote to break ties between two sites.
The Dynamo formula is weighted voting with every weight equal to one. Weights also explain a common mistake: if reads and writes both consult the same fixed subset of heavy nodes, the system behaves like a small cluster with a large, useless periphery. Assign weights from failure domains, not from hardware size.
Grid quorums and the load trade-off
Arrange n = k x k nodes in a grid and define a quorum as one complete row plus one complete column. Any row meets any column, so any two quorums intersect. A quorum has 2k - 1 nodes, which grows with the square root of n instead of linearly. The table compares it with majority, using the exact availability calculation in the code below.
| Nodes | Majority size | Grid size | Unavailability, p = 0.99 | Unavailability, p = 0.9 |
|---|---|---|---|---|
| 9 (3 x 3) | 5 | 5 | majority 1.2e-8, grid 4.6e-5 | majority 8.9e-4, grid 0.033 |
| 16 (4 x 4) | 9 | 7 | majority 1.2e-12, grid 4.6e-6 | majority 6.1e-5, grid 0.025 |
| 100 (10 x 10) | 51 | 19 | not computed | not computed |
import math, itertools
def majority_availability(n, p):
q = n // 2 + 1
return sum(math.comb(n, k) * p**k * (1 - p)**(n - k) for k in range(q, n + 1))
def grid_availability(k, p):
"""Exact: enumerate all 2^(k*k) up/down states; a quorum needs a live row and a live column."""
total = 0.0
for s in itertools.product((0, 1), repeat=k * k):
prob = math.prod(p if up else 1 - p for up in s)
live_row = any(all(s[r * k + c] for c in range(k)) for r in range(k))
live_col = any(all(s[r * k + c] for r in range(k)) for c in range(k))
total += prob if (live_row and live_col) else 0.0
return totalThe numbers show the trade. At 16 nodes a grid quorum is 7 nodes instead of 9, and at 100 nodes it is 19 instead of 51, so each node carries a much smaller share of the traffic. But a grid quorum needs one entire row and one entire column to be alive, so a handful of well-placed failures blocks it. With nodes at 99% uptime, the 16-node grid is unavailable about 4.6 times in a million, against about one in a trillion for majority. Grids suit read-heavy systems with many nodes where throughput matters more than surviving many simultaneous failures; for small clusters, majority wins.
Flexible Paxos: not every pair must meet
Classic Paxos uses majorities for both of its phases. Howard, Malkhi and Spiegelman showed in 2016 that this is more than necessary. Phase 1, where a new leader learns what earlier leaders may have committed, uses quorums Q1. Phase 2, where the leader replicates a value, uses quorums Q2. Safety needs only that every Q1 intersects every Q2. Two phase-2 quorums need not intersect each other, because only one leader is active per ballot and the phase-1 handover carries the history.
With threshold quorums this becomes |Q1| + |Q2| > n. On nine nodes you can replicate with any 3 and elect with any 7. Steady-state replication, the operation you perform millions of times, now waits for the fastest 3 instead of the fastest 5, and tolerates 6 slow or failed followers. The price is paid at election time: a new leader needs 7 live nodes, so with more than 2 failures the current leader can keep running but no replacement can be elected. Designs that use this idea choose it deliberately, for example replication quorums confined to one zone and election quorums that span zones. Always check both directions of the trade before adopting it.
from itertools import combinations
def intersects(q1_sets, q2_sets):
"""True if every quorum in q1_sets shares a node with every quorum in q2_sets."""
return all(a & b for a in q1_sets for b in q2_sets)
def threshold(nodes, k):
return [frozenset(c) for c in combinations(nodes, k)]
def grid(k):
cells = [[r * k + c for c in range(k)] for r in range(k)]
rows = [set(row) for row in cells]
cols = [set(cells[r][c] for r in range(k)) for c in range(k)]
return [frozenset(r | c) for r in rows for c in cols]
nodes = range(9)
print(intersects(threshold(nodes, 5), threshold(nodes, 5))) # True: majority
print(intersects(threshold(nodes, 4), threshold(nodes, 4))) # False: 4 + 4 is not above 9
print(intersects(grid(3), grid(3))) # True: every row meets every column
print(intersects(threshold(nodes, 7), threshold(nodes, 3))) # True: Flexible Paxos Q1 x Q2
print(intersects(threshold(nodes, 3), threshold(nodes, 3))) # False, and Flexible Paxos does not need it
Byzantine quorums
Everything above assumes failed nodes stop. If up to f nodes may lie, returning forged or stale values, one shared node is not enough, because the shared node might be one of the liars. How much more is needed depends on whether data can be verified. Malkhi and Reiter distinguished two cases in 1998. When values are self-verifying, for example signed by the writer, a liar cannot forge them, so it is enough that every intersection contains one honest node: any two quorums must share at least f + 1 nodes. When values are not verifiable and a reader must outvote the liars, the intersection must hold 2f + 1 nodes, which needs n of at least 4f + 1.
Practical protocols such as PBFT authenticate every message and use n = 3f + 1 nodes with quorums of 2f + 1. Check the arithmetic: two quorums of 2f + 1 drawn from 3f + 1 nodes share at least 2(2f + 1) - (3f + 1) = f + 1 nodes. With f = 1, that is four nodes, quorums of three, and every two quorums sharing two nodes, of which at most one can be faulty. Tolerating one liar costs four machines where tolerating one crash costs three, and the saving over 4f + 1 exists only because signatures or MACs let honest nodes prove what they saw.
What breaks intersection in practice
- Sloppy quorums. Writing to substitute nodes during a failure keeps the system available but the substitutes are outside the configured quorum, so a read quorum can miss them. Hinted handoff repairs this later; until then reads can be stale. See hinted handoff.
- Partial writes. A write that reached fewer than W replicas and then failed may still be on some of them. A later read can see it or miss it depending on which replicas answer, so it can appear to flicker. Protocols with read repair make the outcome converge; see read repair.
- Unsafe reconfiguration. Changing membership from an old set of nodes to a new one in one step can leave a moment where a majority of the old set and a majority of the new set do not intersect, so two leaders are elected. Raft offers two fixes. Joint consensus makes decisions during a change need quorums in both the old and new configurations. Single-server changes are safe for a different reason: when the two sets differ by one member, any majority of the old set overlaps any majority of the new set.
- Correlated failures. Availability numbers assume independent failures. Three replicas in one rack are one failure domain. Place quorum members across racks and zones so that a single event cannot take out a quorum's worth of nodes.
- Clock-based shortcuts. Leader leases that let a leader answer reads without a quorum are only safe if clocks are bounded; a paused leader can serve stale reads after a new one is elected.
Choosing a quorum system
| Situation | Choice | Why |
|---|---|---|
| 3-7 node consensus group | Majority | Best availability; load is not the bottleneck at this size |
| Leaderless key-value store | Threshold R and W over N replicas | Tunable per request; see the R, W and N guide |
| Two sites plus a tie-breaker | Weighted votes with a witness | Survives the loss of either site without split brain |
| Many replicas, read-heavy | Grid or other low-load system | Quorum size grows with the square root of n |
| Frequent writes, rare elections | Flexible Paxos quorums | Small replication quorum, large election quorum |
| Untrusted or independently operated nodes | Byzantine quorums, n = 3f + 1 | Intersections must contain an honest node |
Consensus protocols put these quorums to work; Paxos shows the two phases whose quorums Flexible Paxos separates.
What to do next
- List every operation your system performs with a quorum, and for each write down which other operations must see its result.
- Run the checker above against your actual quorum definitions, including during membership changes.
- Compute availability for your node count and a realistic per-node uptime, and compare majority with any alternative you are considering.
- Map quorum members onto racks and zones, and verify that no single failure domain holds enough nodes to block or split a quorum.
- If you use sloppy quorums or partial writes, document the staleness they allow and confirm that repair runs.
- If your write path is latency-bound on a majority, evaluate Flexible Paxos quorums, and test election with failures before adopting them.