Unit content
Data partitioning and sharding
Partitioning divides a dataset so different nodes are responsible for different subsets. In databases this is often called sharding.
A partitioning function maps each key to a shard. Common strategies include ranges, hashes and explicit directory-based assignments.
Partitioning can increase storage and throughput because requests for different shards can be handled independently. It also introduces new problems:
- a request may need data from several shards;
- one hot key or range can overload a single shard;
- adding or removing nodes may require data movement;
- transactions across shards require coordination that one-shard operations can avoid.
Replication and partitioning solve different problems. Replication creates copies of the same data; partitioning distributes different portions of the data.