Most engineers can recite CAP: consistency, availability, partition tolerance, pick two. Fewer can say what each word means in the theorem, and fewer still notice that CAP is silent about the situation a database is in almost all of the time, which is when the network is fine. PACELC fills that gap. It says that if there is a partition (P) a system trades availability (A) against consistency (C), else (E) it trades latency (L) against consistency (C).

This page states CAP precisely, clears up three misreadings, then spends most of its length on the else clause: why replication makes it unavoidable, how to put a number on it, how to classify real systems without overclaiming, and how to choose per operation. The broad survey of CP and AP systems lives in the CAP theorem overview; this page assumes you have seen it or can live without it.

Advertisement

What CAP actually states

Eric Brewer stated CAP as a conjecture in his PODC keynote in 2000. Seth Gilbert and Nancy Lynch proved a formal version in 2002, and the proof is only as strong as their definitions. Consistency means linearizability: every operation appears to take effect at a single instant between its start and its end, so a read never returns a value older than a write that finished before the read began. Linearizability versus sequential consistency covers that model in detail. Availability means every request received by a non-failing node must eventually get a response that is not an error. Partition tolerance means the system keeps to its guarantees even when the network loses an arbitrary number of messages between nodes.

The proof is short. Put two replicas, G1 and G2, on opposite sides of a partition. A client writes v1 to G1, which cannot tell G2. Another client then reads from G2. If G2 answers, it can only return the old value, which breaks linearizability. If G2 waits for the partition to heal, it may wait forever, which breaks availability. No algorithm escapes this, because G2 has no information about the write.

Three misreadings that cause bad designs

  1. Pick any two. You cannot pick CA for a system that runs over a network, because you do not get to choose whether partitions happen. CAP describes what a replicated system does during a partition. A single-node database is CA only in the trivial sense that it has nothing to partition.
  2. CAP availability is uptime. It is not. CAP availability puts no bound on response time, so an answer after an hour counts. It also requires that every non-failing node answers. A Raft cluster where the majority side keeps serving is CP in CAP terms, yet clients routed to the majority see a perfectly available service. The operational word and the theorem word measure different things.
  3. A system is CP or AP. Systems make the choice per operation, and many let the caller choose. Brewer made this point himself in his 2012 essay CAP Twelve Years Later: partitions are rare, so the interesting design work is in detecting them, limiting what runs during them, and recovering afterwards.
Advertisement

PACELC: the else clause

Daniel Abadi proposed PACELC in a 2010 blog post and set it out in a 2012 IEEE Computer paper, Consistency Tradeoffs in Modern Distributed Database System Design. His argument is that the trade-off between consistency and latency shaped the design of real systems more than CAP did, because it applies to every request, not only during a partition.

The trade-off comes from replication itself. Once a write has to reach several machines, there are only a few ways to process it. You can make the replicas agree before acknowledging, which costs at least one round trip to some of them. You can send every write through one primary and acknowledge when the primary has it, which is fast near the primary and slow far from it, and leaves the other replicas briefly behind. Or you can accept the write at any replica and spread it afterwards, which is fastest and means readers elsewhere can see old data or conflicting versions. Each choice moves you along the same axis: wait for coordination, or answer before it finishes.

A request arrivesreplicated dataIs there a partition?replicas cannot talkyes: Pno: E (else)PA: answer anywaymay be stale or divergePC: refuse or waitminority side errorsEL: answer fastack before all replicas agreeEC: coordinate firstpay replica round tripsrare: minutes per yearevery request, all dayA system is classified by the pair it picks, e.g. PA/EL or PC/EC,and a tunable system can pick a different pair for each operation.
PACELC as a decision: the partition branch is the one CAP describes; the else branch is the one your latency dashboards measure.

Classifying systems, carefully

Abadi's paper classified several systems as they behaved by default at the time. These labels are often quoted without his caveats, so here they are with them.

System (2012 default)ClassReasoning
Dynamo, Cassandra, RiakPA/ELGive up consistency for availability under partition, and for latency otherwise. He notes that even with R + W > N they cannot reach Gilbert and Lynch consistency.
VoltDB/H-Store, MegastorePC/ECFully ACID; pay availability and latency to keep consistency.
BigTable, HBasePC/ECSingle owner per row range; consistency first.
MongoDBPA/ECReads go to the primary in the baseline case; a partition causes more consistency issues than availability issues.
PNUTSPC/ELGives up consistency for latency normally; under partition it trades availability, which he admits reads oddly.

Treat the table as a historical snapshot, not a current product sheet. Defaults, replication protocols and client options have changed in most of these systems since 2012. The lasting lesson is the method: ask what the system does to one specific operation in each branch, under the configuration you actually run. For Cassandra, Cassandra's CAP trade-offs in practice does that per consistency level.

Worked example: what consistency costs in milliseconds

Suppose a coordinator in us-east writes to three replicas: one in its own region, one in us-west, one in eu-west, with round trips of roughly 1, 65 and 75 ms. A write acknowledged after k replicas confirm costs the k-th fastest round trip. Latency is an order statistic, so the question is never what consistency costs on average but which replica you are waiting for. The simulation below adds random jitter to each round trip and measures the acknowledgement time for k = 1, 2 and 3.

import random

# One coordinator in us-east. Replica round trips in ms: same region, us-west, eu-west.
BASE = [1.0, 65.0, 75.0]

def ack_latency(k, rng):
    """Time until k of the 3 replicas have acknowledged one write."""
    rtts = sorted(b + rng.expovariate(1 / (0.1 * b + 0.5)) for b in BASE)
    return rtts[k - 1]

def pct(xs, q):
    xs = sorted(xs)
    return xs[int(q * (len(xs) - 1))]

rng = random.Random(7)
for k, name in ((1, "W=1 (ONE)"), (2, "W=2 (QUORUM)"), (3, "W=3 (ALL)")):
    xs = [ack_latency(k, rng) for _ in range(100_000)]
    print(f"{name:13} p50={pct(xs, 0.50):6.1f} ms  p99={pct(xs, 0.99):6.1f} ms")

Running it prints:

W=1 (ONE)     p50=   1.4 ms  p99=   3.7 ms
W=2 (QUORUM)  p50=  69.9 ms  p99=  86.8 ms
W=3 (ALL)     p50=  81.5 ms  p99= 113.2 ms

Waiting for one replica costs a local round trip. Waiting for a majority costs a cross-country one, because the second-fastest replica is always remote in this layout. Waiting for all three adds the slowest replica and its tail: the p99 grows by about 26 ms over quorum. That is the EC price for one write, paid on every request, with no partition anywhere.

Two design levers follow directly. Placement: put a majority of replicas in one region and quorum becomes local, at the cost of losing that region taking the majority with it. Scope: a quorum inside a region (Cassandra's LOCAL_QUORUM) gives consistency for clients in that region and leaves cross-region replication asynchronous, which is EC locally and EL globally. Quorum reads and writes in depth works through the overlap rules behind these levels.

The other side: how stale is EL?

Choosing EL means reads can miss recent writes. How often depends on how soon after a write the read arrives and how long replication takes. The next simulation writes with W=1, then reads one random replica with R=1 after a delay, while the two other replicas apply the write after an exponential lag with a 40 ms mean.

import random

def stale_read_rate(lag_mean_ms, gap_ms, trials, rng):
    """W=1, R=1, N=3: a read goes to a random replica gap_ms after the write.
    The write is on the coordinator replica at once; the other two get it after
    an exponential replication lag. Count reads that miss the write."""
    stale = 0
    for _ in range(trials):
        replica = rng.randrange(3)
        if replica == 0:
            continue
        if rng.expovariate(1 / lag_mean_ms) > gap_ms:
            stale += 1
    return stale / trials

rng = random.Random(11)
for gap in (1, 10, 50, 200):
    print(f"read {gap:3} ms after write: {stale_read_rate(40, gap, 200_000, rng):5.1%} stale")

Running it prints:

read   1 ms after write: 65.1% stale
read  10 ms after write: 52.0% stale
read  50 ms after write: 18.9% stale
read 200 ms after write:  0.4% stale

Immediately after the write, almost two thirds of reads are stale: the ceiling is two thirds, the chance of picking a replica other than the coordinator. By 200 ms the rate falls below one percent. That is the shape to expect, and why EL feels fine for browsing and fails badly for read-your-own-writes flows such as saving a profile and reloading the page. The usual fixes are cheap: route a user's reads to the replica that took their write for a short window, carry a version token and read from a replica that has reached it, or read at quorum only on the page after a write.

Choosing per operation

The useful unit of PACELC is the operation, not the product. An online shop on one replicated store can sensibly run several combinations at once:

-- Same keyspace, three operations, three PACELC choices (Cassandra cqlsh syntax).
CONSISTENCY LOCAL_ONE;      -- browse catalogue: EL, stale by a few ms is fine
SELECT name, price FROM shop.products WHERE id = 42;

CONSISTENCY LOCAL_QUORUM;   -- add to cart: EC inside the region, EL across regions
UPDATE shop.carts SET items = items + {42: 1} WHERE user_id = 7;

-- last unit of stock: a compare-and-set through Paxos, PC/EC for this row only
UPDATE shop.stock SET units = 0 WHERE id = 42 IF units = 1;

The catalogue read tolerates staleness and must be fast, so it is EL. The cart update wants your own reads to see it, so it is EC within the region. The last unit of stock must never be sold twice, so it uses a compare-and-set that goes through a consensus round and refuses to proceed without a majority: PC/EC for that one row, at several round trips of cost. Write down which branch each endpoint takes and why, next to the endpoint, because the choice is invisible in the code.

Consensus systems make the same choice in a different place. A Raft leader can serve reads from local state quickly, but a deposed leader that has not yet noticed can return stale data. Linearizable reads need a round of heartbeats to confirm leadership (read index) or a time-bounded lease. That is EL versus EC inside a CP system; Raft in depth covers both mechanisms.

Failure modes

  • Treating R + W > N as linearizable. Sloppy quorums, hinted handoff, last-write-wins with skewed clocks and concurrent writes all break it, which is Abadi's point about Dynamo-style systems. Overlap gives you a recent value, not a linearizable history.
  • Timeouts treated as answers. A timed-out write may have succeeded on some replicas. Retrying a non-idempotent write after a timeout duplicates it; use idempotency keys or conditional writes.
  • Slow looks like partitioned. A node in a long GC pause is indistinguishable from an unreachable one. Your timeout decides which branch of PACELC you are in, so a tight timeout turns latency spikes into partition behaviour.
  • Partial partitions. Real partitions are often asymmetric: A reaches B, B cannot reach C, the client reaches everyone. Designs that assume a clean split misbehave, for example with two nodes that each believe they lead.
  • Unmeasured staleness. Teams pick EL and never measure replication lag, so the inconsistency window is a guess until a customer finds it.

Operating it

Measure both branches. For the else branch, track p50 and p99 latency by consistency level and replication lag per replica, and alert when lag exceeds the window your read-your-writes trick assumes. For the partition branch, rehearse: inject partitions between racks and regions in staging, including one-way packet loss, and record what each endpoint returns. Fault-injection test harnesses in the style of Jepsen show quickly whether a claimed guarantee holds. Finally, record client timeouts next to consistency levels in configuration review, since the two together decide your real behaviour.

Dynamo-style systems took the PA/EL corner deliberately; the Dynamo paper and its influence explains why that was right for a shopping cart and where the design pushed complexity onto applications.

What to do next

  1. List your top ten endpoints and label each with its PACELC pair, under the configuration you run today.
  2. For every endpoint labelled EC, find the replica it waits for and measure p99 latency for that consistency level.
  3. For every endpoint labelled EL, measure replication lag and decide whether read-your-own-writes needs a session-pinned replica or version token.
  4. Check that every retried write is idempotent or conditional.
  5. Review client and server timeouts together, and widen any that turn routine latency spikes into failover.
  6. Run one partition drill in staging, including a one-way partition, and write down what each endpoint returned.
Key takeaway: CAP says that during a partition a replicated system must choose between linearizable answers and answering at all; it says nothing about the normal case. PACELC adds the normal case: else, choose between low latency and consistency, because coordination costs the round trip to the k-th fastest replica on every request. Classify operations, not products, measure what each choice costs in milliseconds and in staleness, and treat timeouts, clocks and sloppy quorums as the places where claimed guarantees quietly fail.