skip to content

questions

6

In plain terms, what does the CAP theorem say a distributed database must sacrifice when a network partition occurs, and why can't a system just have all three of consistency, availability, and partition tolerance?

level: juniorimportance: must knowfreq 80%

answer

  1. pick 2 of 3, but P isn't optional
  2. partition forces C xor A
  3. split-brain = both sides accept writes
  4. CP = correctness over uptime
  5. AP = uptime over correctness

basics

~20 s

When two halves of a system can't talk to each other, you must pick: either every request works but might return old/wrong data (available), or requests wait/fail until the data is provably correct (consistent). You can't have perfect answers and instant answers at once during that outage.

solid answer

~40 s

CAP theorem states that during a network partition, a distributed system must choose between Consistency (every node returns the most recent write) and Availability (every request gets a non-error response), because the two nodes cut off from each other cannot both stay live and agree on the current value without communicating. Partition tolerance isn't really a design choice — real networks fail, so any system that spans multiple nodes must handle partitions somehow. In practice, CAP reduces to a CP-vs-AP decision: does the system refuse or queue requests on the minority side of a partition to guarantee correctness, or does it keep answering using local state and reconcile divergence later? The choice is usually made per subsystem or per operation, not as a single label for the whole system.

go deeper

for a junior

Should state the basic trade-off in plain language: during a network split, pick correct-but-maybe-unavailable or available-but-maybe-stale. Doesn't need to name Gilbert and Lynch or discuss quorum internals.

for a middle

Should connect the trade-off to a concrete request path — what a node does when it can't reach its peer — and know that P is mandatory in practice, so it's really a CP-vs-AP call.

for a senior

Should discuss failure modes like split-brain and false-positive partition detection, and be able to justify CP vs AP for a specific system with the actual cost named on each side.

for a principal

Should be able to critique CAP's coarseness — it models a single boolean partition and a single all-or-nothing consistency level — and connect it to how real systems tune the trade-off per operation or per key rather than as one global label.

## What the theorem states CAP theorem, first stated by Eric Brewer in 2000 and formally proven by Seth Gilbert and Nancy Lynch in 2002, describes a hard constraint on any system that stores data on more than one machine and talks to those machines over a network. The three properties are: - **Consistency** — every read returns the value of the most recent acknowledged write, as if there were only one copy of the data. - **Availability** — every request that reaches a non-failed node receives a response, not an error or timeout. - **Partition tolerance** — the system continues operating even when messages between nodes are dropped or delayed. The theorem's claim is narrow but important: when a partition actually happens — some nodes can't exchange messages with others — a system can guarantee at most one of Consistency and Availability for the affected data, not both. ## The mechanism The mechanism is concrete. Imagine two replicas, `A` and `B`, holding a copy of the same value, and a network fault severs the link between them. A client on A's side writes a new value. `B` never sees that write because the link is down. Now a second client asks `B` to read the value. `B` has two options: 1. It can answer immediately with the value it has locally — this satisfies **Availability**, but the answer is stale, so Consistency is broken. 2. Or `B` can refuse to answer until it can confirm with `A` that it has the latest data — this preserves **Consistency**, but the request effectively failed, so Availability is broken. There is no third option that lets `B` both answer immediately and guarantee correctness, because correctness requires information `B` does not have and cannot get while the partition persists. ## Why the constraint exists CAP exists because early distributed-systems designers wanted the same guarantees a single-machine database gives — always-consistent reads — while also wanting the scalability and fault tolerance that come from spreading data across many machines and, often, many data centers. Those two goals collide the moment the network misbehaves, and networks always eventually misbehave: - switches fail; - links get saturated; - cross-region links go down for maintenance; - or a node is simply too slow to be distinguished from an unreachable one. Brewer's contribution was naming the trade-off explicitly so architects stop assuming they can dodge it. ## The cost on each side The trade-off has a cost on each side. | The choice during a partition | What it costs and what it buys | |---|---| | **Systems that pick Consistency over Availability** — commonly called CP systems, such as a strongly consistent configuration store or a relational database configured for synchronous replication | They will reject writes, and sometimes reads, on the side of the partition that cannot confirm a quorum. Users on that side see errors or timeouts until the network heals. The benefit is that no client is ever shown incorrect data. | | **Systems that pick Availability** — AP systems, such as many wide-column or key-value stores in their default configuration | They keep serving reads and writes locally on both sides of the partition. The benefit is uptime: no user-facing errors caused by the network fault. The cost is that the two sides can now diverge, and the system needs a reconciliation strategy once the partition heals. | ## Failure modes 1. **Split-brain** is the failure mode most often traced back to CAP in production: both halves of a partitioned cluster believe they are the sole authority and keep accepting writes, so each half's data silently drifts from the other's. This is especially dangerous for systems doing leader election or distributed locking, where two 'leaders' active at once can corrupt shared state. 2. A second common failure mode is a **false partition**: a node isn't actually cut off from the network, but it is so slow — say, during a long garbage-collection pause — that the rest of the cluster times out and treats it as partitioned, triggering an unnecessary failover or a burst of stale reads. 3. A third is teams believing they've built a system with no cost, when in reality they've never tested what happens when the network actually splits, and the day it does, one guarantee quietly breaks. ## Where the choice shows up A concrete illustration: a multi-region deployment of a distributed configuration service like ZooKeeper is deliberately CP. If the network between two data centers is severed, the minority side of the ZooKeeper ensemble stops serving writes and stale reads entirely rather than risk two nodes disagreeing about, say, which service instance currently holds a distributed lock — an incorrect answer there could mean two processes both believe they own an exclusive resource. The minority side becomes unavailable until it can rejoin a quorum, which is the deliberate, accepted cost of guaranteeing correctness for that kind of coordination data.

  • Is it possible for a system to be CA — consistent and available — in a real deployment?
    Only in the degenerate case where the network never partitions, such as a single-node database or a system that treats any network fault as a total outage rather than a partial one. Once a system spans more than one node and must remain reachable across a real, imperfect network, a partition will eventually occur, and the CA label collapses into either CP or AP depending on what actually happens at that moment. Most 'CA' claims in practice mean the team hasn't yet observed or tested a partition.
  • How do quorum-based systems decide which side of a partition stays available?
    The side that can still reach a majority of the total nodes keeps serving reads and writes, since a majority guarantees it cannot conflict with any other majority — two disjoint majorities of the same node set can't both exist. The minority side, unable to form a quorum, stops accepting writes and typically stops serving reads too, to avoid returning stale data. This is why quorum systems are usually deployed with an odd number of nodes across independent failure domains.
  • Does choosing availability during a partition mean the system has no consistency guarantees at all?
    No — AP systems typically still offer some consistency guarantee, just a weaker one than always-latest, and reconcile divergent copies once the partition heals using mechanisms like timestamps or version vectors. The exact strength of that weaker guarantee is a separate, deeper topic than CAP itself; CAP only tells you that during the partition, the strongest guarantee isn't available.

Two bank branches lose their phone line to each other. If a customer withdraws the last $100 from Branch A, Branch B either has to refuse withdrawals on that account until the line is back up (consistent but unavailable), or let a customer withdraw from B too, risking the account going negative because B didn't know about A's withdrawal (available but inconsistent).

saying these in an interview costs you the question

  • Claims a system can be genuinely 'CA' in a real multi-node network deployment
  • Says CAP means permanently picking 2 of 3 for the whole system rather than a per-partition, often per-operation choice
  • Can't explain what actually happens to a request on the minority side of a partition
  • Treats 'eventual consistency' and 'no consistency at all' as the same thing
  • Doesn't mention that partition tolerance isn't optional for a real networked system

context

open as a page

PACELC extends the CAP theorem with a clause for when the network is NOT partitioned. What does the 'ELC' part of PACELC add, and how does it change a system's design compared to only thinking in CAP terms?

level: middleimportance: must knowfreq 55%

basics

~20 s

PACELC says: if there's a network split (P), pick availability or consistency (A or C) — that's CAP. But Else (E), even when the network is fine, you still must pick between faster answers (Latency) or more correct answers (Consistency), because waiting to check with other machines always costs time.

open as a page

Some engineers describe a system as 'CA' — consistent and available, with no partition tolerance needed. Why do most distributed-systems practitioners consider 'CA' not a real option for a system that operates over an actual network, and what is the person usually actually describing when they use that label?

level: middleimportance: must knowfreq 45%

basics

~20 s

Real networks break sometimes — cables get cut, switches fail, regions lose connection. A system that spans more than one machine over such a network WILL face a partition eventually, so it can't promise to just never deal with one. 'CA' usually really means 'we haven't seen a partition happen yet' or 'we treat the whole system as down if one does.'

open as a page

You're picking a data store for a distributed lock service used for leader election versus a data store for shopping-cart sessions on a high-traffic e-commerce site. Walk through why you'd lean CP for one and AP for the other, and name the concrete cost each choice imposes on that specific system during a network partition.

level: seniorimportance: must knowfreq 70%

basics

~20 s

For the lock service, two 'leaders' at once could break everything, so it's worth going offline rather than risk that (CP). For shopping carts, showing an old cart is annoying but rarely catastrophic, and going offline during checkout loses sales, so keep answering even with slightly stale data (AP).

open as a page

A wide-column store lets you configure the consistency level per read or write request — for example, requiring acknowledgment from just one replica versus requiring acknowledgment from a majority of replicas. How does this per-request knob let a team make different CAP and PACELC trade-offs for different operations within the SAME cluster, and what's the risk of misusing it?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Instead of the whole database picking 'always fast' or 'always safe' once, each individual read or write can choose: talk to just one copy (fast, might be stale) or talk to most copies (slower, more sure to be correct). Different parts of the app can pick differently. The risk is picking 'fast' for something that actually needed to be 'safe' and getting wrong data where it matters.

open as a page

CAP theorem is frequently criticized by distributed-systems researchers as too coarse a model to directly drive real architecture decisions. What are the main limitations of the CAP formulation itself, and in what way does PACELC address one of them while still leaving others unresolved?

level: principalimportance: nice to knowfreq 25%

basics

~20 s

CAP treats 'consistent' and 'available' as simple yes/no labels for a whole system, but real systems have many shades of both, change behavior operation-by-operation, and CAP says nothing about what happens when the network is fine. PACELC fixes the 'network is fine' gap but still treats things too simply.

open as a page