The CAP theorem states that when a distributed system experiences a network partition, it can guarantee either strong consistency or high availability, but cannot guarantee both. Because network cables disconnect and switches fail, partition tolerance is mandatory for distributed infrastructure, forcing architects to make an explicit trade-off between serving stale data or rejecting writes during faults. Technical interviewers test whether you understand this operational trade-off rather than just reciting the acronym.
The network partition trade-off
In a distributed environment, partition tolerance means the cluster continues operating despite dropped packets or severed connections between nodes. Because physical networks will inevitably experience delays or splits, you cannot select CA in a multi-node deployment. When inter-node communication breaks down, you face two strict paths:
- CP systems: Choose consistency over availability. When a partition isolates nodes, the system halts writes or returns errors on the minority partition to prevent split-brain state divergence. Apache HBase and Google Cloud Spanner follow this model, refusing mutations if quorum or synchronized TrueTime clocks cannot guarantee linearizable order.
- AP systems: Choose availability over consistency. Nodes accept writes and serve reads even when disconnected from peer nodes. Apache Cassandra and Amazon DynamoDB (by default) follow this approach, returning local values immediately and resolving divergence later via background anti-entropy repairs.
Network Partition Occurs: Node A (Majority) <--- Network severed ---> Node B (Minority) CP Decision: Reject writes on Node B to prevent state divergence. AP Decision: Accept writes on both nodes; reconcile conflicting versions later.
Normal operations without partitions require additional latency considerations:
PACELC and day-to-day operations
The PACELC theorem extends CAP by explaining system behaviour during normal operations: if there is a Partition (P), how does the system choose between Availability (A) and Consistency (C); Else (E), how does it trade Latency (L) versus Consistency (C)?
Under normal network conditions, providing strong consistency requires inter-node communication, synchronous round-trips, and quorum consensus, which adds query latency. Choosing low latency means writing locally and returning early, accepting eventual consistency. When selecting data stores for pipelines, choose CP stores like Spanner or PostgreSQL replicas with synchronous replication for financial ledgers, and AP stores like Cassandra or DynamoDB for high-throughput clickstream ingestion where write dropouts cannot be tolerated.