/knowledge/notes/cap-sharding-replication
Concept note · ML
CAP, Sharding and Replication
Distributed Systems
- Studied
- Cluster and Cloud ComputingCOMP90024
- When
- 2023 S1
- Applied in
- Social Sense
- Read / Refreshed
- ~5 min read2026-10-15
The CAP theorem says a distributed data store can only guarantee two of three properties: Consistency, Availability, and Partition tolerance. In practice, network partitions are inevitable, so you choose between CP (consistent but unavailable during partitions) and AP (available but eventually consistent). Sharding and replication are the mechanisms that implement these choices.
01
The idea
Consistency means every read sees the most recent write. Availability means every request gets a response (success or failure). Partition tolerance means the system works even when network failures split it into isolated groups. The CAP theorem says you cannot have all three when a partition occurs.
Replication copies data across nodes so reads can be served locally and writes survive failures. Sharding splits data across nodes so the system scales horizontally. Together they determine how your system behaves during failures.
In a CP system, when a partition happens, nodes in the minority partition refuse to serve requests to avoid returning stale data. This sacrifices availability for consistency. In an AP system, all partitions keep serving requests, but they might return different answers until the partition heals and they reconcile.
Most real systems are neither pure CP nor pure AP. They offer tunable consistency: you can choose per-request whether to wait for acknowledgment from all replicas (strong consistency, low availability) or just one (high availability, eventual consistency).
02
The maths
In a replicated system with n nodes and replication factor r, each data item is stored on r nodes. A read quorum of q_r nodes and a write quorum of q_w nodes must satisfy q_r + q_w > r to guarantee consistency, because any read overlaps with the latest write.
For strong consistency (linearizability), you need q_r + q_w > r and q_w > r/2. The common choice is q_r = q_w = ⌈(r+1)/2⌉, a majority quorum. If r = 3, you need 2 nodes to agree. This tolerates ⌊r/2⌋ node failures while maintaining consistency.
In sharded systems, data is partitioned across s shards using a hash function or range partitioning. The probability that a random key is on a failed shard is f/s, where f is the number of failed shards. Replicating each shard r times reduces the failure probability to (f/s)^r, assuming independent failures.
The latency of a quorum read is the median of response times from q_r replicas. The tail latency matters: if one replica is slow, the 99th percentile request latency is often dominated by the slowest replica. Hedging (sending duplicate requests) reduces tail latency at the cost of increased load.
03
Try it
The widget simulates a three-node cluster with a network partition. In CP mode, the minority partition refuses writes to stay consistent. In AP mode, all partitions accept writes, creating divergent state. When the partition heals, AP systems must reconcile conflicting writes using vector clocks or last-write-wins.
04
Where I used it
05
Easy to get wrong
06
Sources
Covered in COMP90024 Cluster and Cloud Computing. Based on Brewer's CAP theorem (2000), Gilbert & Lynch's formal proof (2002), and Kleppmann's Designing Data-Intensive Applications (2017), Chapter 9. Dynamo and Cassandra papers are canonical AP examples; Google Spanner is CP.