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.

Advertisement

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.

Advertisement

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.

Three quorum systems over nine nodesMajority: any 5 of 9nodes 0-4 and 4-8 share node 4Grid 3 x 3: a row plus a columnrow 1 plus column 2 = 5 nodesFlexible Paxos, n = 9Phase 1 (election)any 7 nodesPhase 2 (replication)any 3 nodes7 + 3 is greater than 9The only rule that matters: the quorums that must see each other's work intersectwhich pairs must intersect, and how big they are, is the whole design spaceMajority maximises availability for a given n. Grids lower the load per node. Flexible Paxos makes the commonoperation cheap and the rare one expensive. Byzantine systems need every pair to share at least f + 1 nodes.
Majority, grid and Flexible Paxos quorums over nine nodes. Each is a different answer to which quorums must intersect.
NodesMajority sizeGrid sizeUnavailability, p = 0.99Unavailability, p = 0.9
9 (3 x 3)55majority 1.2e-8, grid 4.6e-5majority 8.9e-4, grid 0.033
16 (4 x 4)97majority 1.2e-12, grid 4.6e-6majority 6.1e-5, grid 0.025
100 (10 x 10)5119not computednot 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 total

The 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

SituationChoiceWhy
3-7 node consensus groupMajorityBest availability; load is not the bottleneck at this size
Leaderless key-value storeThreshold R and W over N replicasTunable per request; see the R, W and N guide
Two sites plus a tie-breakerWeighted votes with a witnessSurvives the loss of either site without split brain
Many replicas, read-heavyGrid or other low-load systemQuorum size grows with the square root of n
Frequent writes, rare electionsFlexible Paxos quorumsSmall replication quorum, large election quorum
Untrusted or independently operated nodesByzantine quorums, n = 3f + 1Intersections 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

  1. List every operation your system performs with a quorum, and for each write down which other operations must see its result.
  2. Run the checker above against your actual quorum definitions, including during membership changes.
  3. Compute availability for your node count and a realistic per-node uptime, and compare majority with any alternative you are considering.
  4. Map quorum members onto racks and zones, and verify that no single failure domain holds enough nodes to block or split a quorum.
  5. If you use sloppy quorums or partial writes, document the staleness they allow and confirm that repair runs.
  6. If your write path is latency-bound on a majority, evaluate Flexible Paxos quorums, and test election with failures before adopting them.
Key takeaway: A quorum system is defined by one property: every quorum that must observe an operation intersects every quorum that performs it. Majorities give the best availability and are right for small consensus groups. Weighted votes express unequal failure domains, and grids trade availability for much lower load. Flexible Paxos notices that only election and replication quorums must meet, so replication can be cheap and elections expensive. With authenticated messages, Byzantine quorums need f + 1 shared nodes and 3f + 1 machines. Sloppy quorums, careless reconfiguration and correlated failures quietly break intersection, so check them explicitly.