← Database Concepts

Beyond Relational Databases: The CAP Theorem and Cassandra's Trade-offs

Published on 2026-08-20·v1.0

Objective

Get oriented before going deep. Everything else in the Cassandra category — the CQL data model, quorum math, write and read paths, gossip and vnodes — assumes you already believe a shared-nothing, eventually-consistent database is worth the trouble. This concept is the argument for that belief: why relational databases hit a wall at web scale, what Brewer's CAP theorem actually says you must give up when a network partition happens, and where Cassandra deliberately lands on that spectrum. cassandra-consistency-levels picks up immediately after this one — it explains the mechanics of how Cassandra lets you dial consistency up or down per query; this concept explains why that dial has to exist at all.

Jeff Carpenter and Eben Hewitt open Cassandra: The Definitive Guide by asking "What's Wrong with Relational Databases?" and then, deliberately, answering their own question: "So the short answer... is 'Nothing.'" The honest version of the story is not that relational databases are bad — it's that every database is a set of trade-offs suited to certain problems, and the trade-offs that make relational databases excellent at consistency and rich querying are the same ones that make them expensive to scale horizontally. CAP theorem is the formal statement of why you cannot dodge that trade-off by trying harder; Cassandra's whole design is a specific, deliberate answer to it.

Use Cases

  • Defending a choice of Cassandra (or push-back on that choice) in an architecture review, using the actual scaling mechanics — locking, joins, two-phase commit — rather than a vague "it doesn't scale" about the relational alternative.
  • Onboarding a developer coming from a single-node ACID database who keeps asking "but what happens if the write only reaches one replica?" — CAP theorem is the vocabulary for answering precisely.
  • Deciding, for a specific workload, whether Cassandra's AP-leaning default is actually the right fit, versus a CP system like a relationally-sharded Postgres cluster or a NewSQL database like CockroachDB or Spanner.
  • Explaining to a team why "Cassandra is eventually consistent" is not the whole story — and pointing them at cassandra-consistency-levels for the mechanism that lets a single keyspace serve both a QUORUM account balance and a ONE telemetry write.
  • Recognizing sharding, shared-nothing architecture, and denormalization as the same set of ideas whether you meet them as hand-rolled patches on a relational database or as Cassandra's default architecture — because the second is the industrialized, automated version of the first.
  • Catching a design decision that quietly buys availability at the cost of a correctness guarantee nobody signed off on losing — the frequent bug is not "we chose AP," it's "we chose AP without discussing it."

Deep Dive

What actually breaks when a relational database scales

The book's argument is concrete, not abstract. Relational consistency is enforced through transactions, and "the way that databases gain consistency is typically through the use of transactions, which require locking some portion of the database so it's not available to other clients. This can become untenable under very heavy loads, as the locks mean that competing users start queuing up, waiting for their turn to read or write the data."

The standard escalation path — the book walks through it step by step — is vertical scaling (bigger box), then a database cluster (which introduces replication and consistency problems you didn't have before), then configuration tuning (turning off logging, which is often not legally an option), then query and index optimization, then a caching layer (which introduces its own cache-versus-database consistency problem), and finally denormalization, which is "antithetical to the five normal forms that characterize the relational model." Each step buys time; none of them removes the underlying tension.

Distributing a relational database's transactions across multiple machines means coordinating them with a two-phase commit. The book is blunt about the cost: "Two-phase commit blocks; that is, clients... must wait for a prior transaction to finish before they can access the blocked resource," and "the problems that 2PC introduces for application developers include loss of availability and higher latency during partial failures." A protocol built to preserve strict consistency across nodes turns out to cost you availability precisely when a node is unreachable — which is CAP theorem showing up in practice before the book ever names it.

Sharding is the other lever, and the book frames Cassandra's architecture as its logical endpoint: "Cassandra uses an approach similar to key-based sharding to distribute data across nodes, but does so automatically." More fundamentally, Cassandra is a shared-nothing architecture — "there is no centralized (shared) state, but each node in a distributed system is independent, so there is no client contention for shared resources... The Cassandra database is a shared-nothing architecture, as it has no central controller and no notion of primary/secondary replicas; all of its nodes are the same." Every node being interchangeable is what makes horizontal scaling close to linear — and it's also what makes "does everyone agree on the current value?" a genuinely hard question the moment the network hiccups.

Brewer's CAP theorem: the trade-off made formal

Eric Brewer stated the conjecture at PODC 2000; Seth Gilbert and Nancy Lynch gave it a formal proof two years later. The claim is that a distributed data system can provide at most two of three guarantees simultaneously:

  • Consistency — every node in the system returns the same, most recent value for a given piece of data; a read never observes a stale or divergent answer.
  • Availability — every request to a non-failing node receives a response, without guaranteeing it's the most recent write.
  • Partition tolerance — the system keeps operating even when network failures split it into groups of nodes that cannot talk to each other.

In a single-datacenter, single-machine relational database, this trade-off is close to invisible — there's nothing to partition. It becomes unavoidable the moment you have more than one node, because a network between machines will fail sometimes, and the theorem's real content is about what the system does at that moment. Once a partition happens, a node has to choose: answer the request anyway, possibly with stale or conflicting data (favoring Availability), or refuse to answer until it can confirm it has the current value (favoring Consistency). Since a real distributed system cannot opt out of network partitions happening eventually, "pick two of three" resolves in practice to picking C or A specifically for the partition case — which is why CAP is usually shorthand as "CP" or "AP," not a genuine three-way menu.

Where Cassandra sits: AP with tunable consistency

Cassandra's own documentation states the choice plainly: "High availability is a priority in web-based applications and to this objective Cassandra chooses Availability and Partition Tolerance from the CAP guarantees, compromising on data Consistency to some extent." That makes Cassandra an AP database by default — a write during a partition succeeds on whatever replicas it can reach rather than blocking, and "eventual consistency is a tradeoff to achieve high availability."

But "AP by default" undersells what Cassandra actually offers, and this is the seam where this concept hands off to cassandra-consistency-levels: Cassandra doesn't force one point on the CAP spectrum for the whole cluster. Consistency level is chosen per query, by the client, against a replication factor set per keyspace — so the same cluster can run one query at ONE (fast, available, eventually consistent) and another at QUORUM or ALL (slower, tilted toward consistency, less tolerant of a down replica) simultaneously. In CAP terms, tunable consistency doesn't cancel the theorem — a single query at ALL during an actual partition still fails, exactly as CAP predicts — it just means "where does this system sit on the C/A line" stops being a single cluster-wide answer and becomes a per-query decision the application makes deliberately (or, worse, accidentally).

Book vs today

PACELC is the extension practitioners now reach for alongside CAP, not instead of it. CAP only describes behavior during a network partition, which the book's own framing hints at without stating outright. Daniel Abadi's 2010 PACELC formulation makes the gap explicit: even when there is no partition, a distributed system still has to trade Latency for Consistency on every request, because confirming a value is current across replicas takes time a purely local read wouldn't. PACELC reads as "if Partition, then Availability vs. Consistency; Else, Latency vs. Consistency" — Cassandra is PA/EL under that scheme: available over consistent during a partition, and latency-favoring over consistency-favoring the rest of the time, which is exactly what choosing ONE or LOCAL_QUORUM as your everyday consistency level buys you. The theorem itself hasn't changed; what's changed is that citing CAP alone, in 2026, reads as citing half the relevant theory.

Brewer's own 2012 retrospective walked back "2 of 3" as a framing, not as a result. In "CAP Twelve Years Later," Brewer wrote that "the '2 of 3' formulation was always misleading" because it suggests a permanent, cluster-wide, binary choice. His corrections track closely onto what tunable consistency already does in practice: "the choice between C and A can occur many times within the same system at very fine granularity... the choice can change according to the operation," and "all three properties are more continuous than binary" — availability is a percentage, not a boolean, and there are "many levels of consistency" between strict and none. Cassandra's per-query consistency level is a working implementation of exactly the fine-grained, continuous choice Brewer says the "2 of 3" slogan obscured. The underlying proof (Gilbert & Lynch) is unchanged; the popular description of it has gotten more precise.

Trade-offs

  • CAP is a statement about partitions, not a permanent cluster identity. Calling Cassandra "an AP database" is a useful shorthand but an oversimplification the theorem's own author disavowed: the real behavior is negotiated per query, per keyspace, and only actually tested the moment a partition occurs. A cluster that has never seen a network partition hasn't validated its CAP choice at all — it's running on the availability and consistency levels its configuration implies, untested.
  • Choosing availability over consistency is not free just because it's the default. The book's escalation story — locks, 2PC, caching layers, denormalization — exists because every one of those relational patches was itself an attempt to buy back some availability or scalability at the cost of some consistency or simplicity. Cassandra doesn't remove that trade, it just moves the decision earlier and makes it explicit at the architecture level instead of implicit in query tuning.
  • Shared-nothing architecture buys near-linear horizontal scaling and gives up cross-node coordination for free. "No centralized state, no primary/secondary replicas" is precisely what lets Cassandra scale by adding nodes; it's also precisely why there's no built-in referee to say authoritatively "this is the current value" the instant two replicas disagree — that job moves to consistency levels, hinted handoff, and repair, each with their own partial coverage (see cassandra-consistency-levels).
  • PACELC's latency-versus-consistency axis applies even on a healthy cluster, not just during outages. Teams that reason about CAP alone often assume "no partition, no trade-off" — but every QUORUM read pays a latency cost relative to ONE even when every node is reachable, because it still waits on the slower of two replicas. Sizing a consistency level means asking both CAP's question (what happens if a replica is unreachable?) and PACELC's (what does this cost on the happy path, forever?).
  • "Eventually consistent" describes a guarantee about the future, not a bound on how long "eventually" takes. The book's own escalation path shows the cost of not facing this squarely: a cache-versus-database consistency problem, once introduced, is exactly the kind of open-ended staleness window that "AP with tunable consistency" makes official rather than accidental — but official doesn't mean bounded. Anti-entropy repair, hinted handoff, and the consistency level chosen are what actually bound that window in Cassandra; CAP theorem alone says nothing about how long convergence takes.
  • A one-size-fits-all CAP choice is usually wrong for a real application. The same argument the book uses against relational databases as a "one-size-fits-all solution" applies to picking a single point on the CAP/PACELC spectrum for an entire system. An account-balance write and a click-stream event have different tolerances for staleness and different tolerances for latency; Cassandra's answer is to let that decision be made per query rather than baked into the database's architecture once, for everyone, forever.

Documentation Links