Your single database is now the bottleneck - too much data, too many writes, too many reads. Two techniques solve two different problems:
Replication (solves READ scaling):
┌──▶ [Read Replica 1]
[Primary DB] ─────┼──▶ [Read Replica 2]
(writes) └──▶ [Read Replica 3]
All writes go to the primary. Reads get spread across replicas, which stay in sync via replication. Great when you have way more reads than writes (true for most apps).
⚠️ Watch for replication lag - replicas can be milliseconds to seconds behind the primary. If a user posts a comment and immediately refreshes, they might not see it yet if they're routed to a lagging replica. This is a classic system design follow-up question.
Sharding (solves WRITE scaling and storage limits):
[Shard 1: users A-H] [Shard 2: users I-P] [Shard 3: users Q-Z]
Split your data across multiple databases, each holding a subset. Now write load AND storage is distributed, not just reads.
The hard part interviewers dig into: how do you pick a shard key? Pick badly (like splitting alphabetically by name) and you get "hot shards" - massively uneven load, since names aren't evenly distributed. A better key is often something like
user_id % number_of_shards, or a hash of the ID, to spread load evenly.The other hard part: cross-shard queries (like "find all users who did X across every shard") become expensive, since you often need to query every shard and merge results.
If you were sharding a system like Instagram by
user_id, what's one query that would suddenly become painful? 👇