Replication copies identical data across multiple physical nodes to provide high availability, fault tolerance, and horizontal read scaling. Sharding partitions distinct subsets of data across separate nodes using a shard key, enabling horizontal write scaling and unbounded storage growth. While replication guarantees survival if an individual server dies, sharding prevents datasets from exceeding the disk and CPU limits of a single machine, and production systems almost always combine both techniques.
Splitting versus duplicating state
The functional distinctions between replication and sharding center on what physical resources each strategy scales:
- Replication mechanics: Every replica node maintains a duplicate copy of the dataset. When the primary node crashes, an automated election promotes a secondary replica to take over traffic without data loss. Read queries can be distributed across read replicas to scale query throughput, but every node must still store the entire database volume and process every write.
- Sharding mechanics: The total dataset is divided into logical partitions called shards. Node A holds customer records from A through M, while Node B holds records from N through Z. Because each node handles only its assigned fraction of total writes and storage, sharding breaks through the physical storage ceilings of single servers.
Replication: Node 1: [Dataset A + B + C] (Primary) Node 2: [Dataset A + B + C] (Read Replica) Node 3: [Dataset A + B + C] (Read Replica) Sharding: Node 1: [Dataset Partition A] Node 2: [Dataset Partition B] Node 3: [Dataset Partition C]
Combining these strategies introduces operational complexities in key selection:
Shard selection and rebalancing risks
Distributed clusters typically shard data into partitions and replicate each partition across three nodes for durability. The primary operational danger in sharded architectures is poor shard key selection:
- Hotspots: Selecting a low-cardinality shard key, such as country code or tenant type, directs disproportionate traffic to a single physical node, causing resource exhaustion while peer nodes sit idle.
- Monotonic keys: Using auto-incrementing IDs or insertion timestamps routes all concurrent writes to the current active shard, eliminating horizontal write distribution.
- Cluster rebalancing: Adding nodes to a sharded cluster requires migrating partitions over the network. Rebalancing consumes disk I/O and network bandwidth, which can degrade query latency if not throttled.