Learning path

Full curriculum

Full curriculum

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.