Distributed Systems Basics: CAP, Consensus, and Eventual Consistency

A program running on one computer lives in a tidy world. Memory is reliable, operations happen in order, and if the machine dies, the whole program dies with it, cleanly. The moment you spread that program across several machines, for scale or for resilience, you leave that tidy world behind. Parts of the system can fail while other parts keep running, messages between machines can be slow or lost, and two nodes can hold different ideas about what is true at the same instant. Distributed systems are the study of building something dependable on top of that mess, and the core ideas are worth knowing even if you never design one from scratch.

Why distribution is genuinely hard

The root difficulty is partial failure. On a single machine, a component either works or the machine is down. In a distributed system, one node can fail while the rest carry on, and worse, the healthy nodes often cannot tell the difference between a peer that has crashed and a peer that is simply slow or unreachable. A request that got no reply might mean the work never happened, or it might mean the work happened and the confirmation got lost. You cannot always tell which.

On top of that, the network itself is unreliable. Messages arrive late, arrive out of order, arrive twice, or never arrive. There is a well known list called the fallacies of distributed computing[1], cataloging the comfortable assumptions that quietly ruin systems, starting with "the network is reliable" and "latency is zero." Every hard problem in the field traces back to the fact that those assumptions are false.

The CAP theorem, stated correctly

The CAP theorem is the most cited and most misquoted idea[2] in the field. It concerns three properties: consistency, meaning every read sees the most recent write, availability, meaning every request gets a non-error response, and partition tolerance, meaning the system keeps working when the network drops messages between nodes.

The popular version says "pick two of three." That phrasing is misleading. In any real distributed system, network partitions are not optional. They will happen, so partition tolerance is not a choice you get to decline. What the theorem really says is narrower and more useful: when a partition is happening, and only then, you must choose between consistency and availability. You can refuse to answer until you are sure your data is current, which sacrifices availability, or you can answer with possibly stale data, which sacrifices consistency. When the network is healthy, you can have both. CAP is a statement about the trade-off forced on you during a failure, not a menu you order from on day one.

Consistency is a spectrum, not a switch

Outside of CAP's strict definition, consistency in practice is a range of models, each a different promise about what a reader can expect to see.

Strong consistency behaves as if there were a single copy of the data[3]. Once a write succeeds, every later read everywhere sees it. This is the easiest model to reason about and the most expensive to provide, because the nodes must coordinate on every operation.

Eventual consistency makes a weaker promise[4]. If writes stop, all copies will converge to the same value, eventually, but for a while different nodes may return different answers. This is what lets a system stay fast and available across many machines and regions. In between sit models like read your writes, which guarantees you at least see your own changes[5] even if others lag, which is often exactly enough for a good user experience.

The practical point is that stronger consistency costs latency and availability, and weaker consistency buys speed and resilience at the price of temporary disagreement. The art is matching the model to the feature. A bank balance wants strong guarantees. A count of likes on a post can be eventually consistent and nobody is harmed if it is briefly off by one.

Consensus: agreeing on a single truth

Many systems need a group of machines to agree on one value or one ordering of events, even though some of those machines may fail and the network may misbehave. That problem is called consensus[6], and it is the beating heart of reliable distributed storage.

The first widely studied solution was Paxos, described by Leslie Lamport[7], which is correct and notoriously hard to understand and implement. Raft came later, designed explicitly to be understandable[8], and it has become the common choice. The rough idea in both is that the nodes elect a leader, the leader proposes each change, and a change is only committed once a majority of nodes have acknowledged it[9]. Requiring a majority, a quorum, is what lets the system tolerate a minority of nodes failing while never committing two conflicting decisions. Systems like etcd, Consul, and the coordination layers inside many databases run Raft[10] to keep a cluster agreeing on its state.

Picking your trade-off on purpose

None of this tells you the right answer, because there is no universally right answer. There is only the trade-off you choose deliberately versus the one that chooses you during an outage.

A system that favors consistency will, during a partition, return errors rather than risk serving stale or conflicting data. Many financial and inventory systems make that choice, because a wrong answer is worse than no answer. A system that favors availability will keep serving during a partition and reconcile differences afterward. Shopping carts, social feeds, and large content platforms often choose this, because being down costs more than being briefly inconsistent. The failure of many systems is not choosing wrong, it is not choosing at all, and discovering the behavior for the first time in the middle of an incident.

The takeaway

Distributed systems are hard for one honest reason: parts can fail independently and the network cannot be trusted, so perfect knowledge is impossible. The CAP theorem names the trade-off a partition forces between consistency and availability. Consistency itself is a spectrum from strong to eventual, and you pick a point on it per feature. Consensus algorithms like Raft are how a cluster agrees on a single truth despite failures, using majority quorums. You do not need to implement any of this by hand to benefit from understanding it, because it tells you what questions to ask of every database, queue, and cache you rely on.

Sources (10)
  1. Wikipedia: Fallacies of distributed computing
  2. Wikipedia: CAP theorem
  3. Wikipedia: Strong consistency
  4. Wikipedia: Eventual consistency
  5. Wikipedia: Consistency model
  6. Wikipedia: Consensus (computer science)
  7. Wikipedia: Paxos (computer science)
  8. In Search of an Understandable Consensus Algorithm
  9. The Raft Consensus Algorithm
  10. The Raft Consensus Algorithm: Implementations