System design2 min

CAP & PACELC Theorems

When building distributed data systems, you cannot have everything. The CAP Theorem proves that a distributed system can only provide two of the following three guarantees simultaneously.

The CAP Theorem

  1. Consistency (C): Every read receives the most recent write or an error. If User A updates their profile, User B must see that update instantly, no matter which node they read from.
  2. Availability (A): Every request receives a (non-error) response, without the guarantee that it contains the most recent write. The system is always up and serving data.
  3. Partition Tolerance (P): The system continues to operate despite an arbitrary number of messages being dropped (or delayed) by the network between nodes. Network partitions are inevitable in the real world.

The Reality of CAP

Because network partitions (P) are a physical reality of distributed systems (cables get cut, switches fail), you cannot choose CA. You are forced to choose between CP and AP.

  • CP (Consistent & Partition Tolerant): If a network link breaks, the system will reject requests (sacrificing Availability) rather than risk returning stale data.
    • Example: Banking systems. Better to show an error than let someone withdraw money they don't have.
  • AP (Available & Partition Tolerant): If a network link breaks, nodes will continue answering queries based on the (potentially stale) data they currently have, sacrificing Consistency.
    • Example: Social media feeds. It's fine if you see a post 5 seconds after your friend sees it; keeping the site online is more important.

The PACELC Theorem

CAP is a great theoretical model, but it only applies when the network breaks (during a partition). What happens during normal, healthy operations? The PACELC theorem extends CAP to cover this.

It states: In case of a Partition, you must choose between Availability and Consistency (CAP). Else (during normal operation), you must choose between Latency and Consistency.

Latency vs Consistency

If your network is healthy, you still have a trade-off:

  • Do you want extremely low Latency? Then you must accept eventual consistency (e.g., asynchronous replication where nodes sync data slightly behind the scenes).
  • Do you want perfect Consistency? Then you must accept higher latency (e.g., synchronous replication where every write forces all nodes to lock and update before returning to the user).

Cassandra, for example, is typically tuned as a PA/EL system (prioritizes Availability during partitions, and Latency during normal operations). Relational databases are usually PC/EC (prioritizes Consistency at all times).