AI System Design
← Learn
Level 3advanced

Database Sharding

One database can't hold it all. Split it.

Depth:
1

Mission

Your dataset and write volume outgrew a single database. Partition the data so it scales.
2

Interactive Simulation

Before we explain anything — play. Push it until it breaks, then fix it.

One database, 6.0K writes/s, 2.5K/s capacity. Replicas won't help with writes — split the data.

Writes per shard

user query touches 1.0 shard
6.0K
S1

Bar height = load vs. each shard's 2.5K/s capacity.

Requests
6.0K/s
Latency
1.21s
p95 3.62s
Error rate
58.3%
CPU
100%
DB load
100%
Est. cost
$420/mo
illustrative
Accepted
2.5K/s
Rejected
3.5K/s
Latency (ms)
Error rate (%)

Load

Partitioning

Strategy
Shard key

Break it

System Score34
3

What just happened?

Sharding split the data across multiple databases by a shard key, so writes and storage spread out — but a bad key created hot partitions.

4

The concept

Sharding (horizontal partitioning) splits data across multiple databases by a shard key. Range sharding splits by key ranges; hash sharding distributes by a hash of the key. Unlike replication, sharding scales writes and storage, not just reads.

5

Trade-offs

Nothing is free. Here's what this solution costs you.

Hot partitions
A skewed shard key overloads one shard.
Cross-shard queries
Joins and transactions across shards are hard and slow.
Salting costs reads
Spreading a hot key over every shard means reading it back fans out.
6

Mini quiz

Question 1 of 30 correct

Sharding scales what that replication does not?

7

Interview me

The app becomes your interviewer. One question, in your own words.

8

Boss challenge

Design a shard key

12K writes/s, and a celebrity account is driving 30% of them.

Goal: Keep errors under 1% and the hottest shard under 90% load.

Use the simulator above with no hints. These checks update live as you play.

9

Interview question

“Explain range vs hash sharding, how to choose a shard key, and how to avoid hot partitions.”

Next: CAP & Consistency