Sharding is horizontal database scaling: data is split into partitions (shards) distributed across separate servers. A single server cannot handle billions of rows — sharding distributes the load.
Partitioning strategies
Range-based — shard 1: id 1–1M, shard 2: id 1M–2M. Simple, but uneven load
Hash-based — shard = hash(id) % N. Even distribution, but painful resharding
Directory-based — a separate catalogue service knows where each record lives. Flexible, but an extra hop
Challenges
JOINs across shards are impossible or very expensive