Replication keeps copies of the same data on several machines, so reads scale and a failure loses nothing; sharding splits the data, so each machine holds only part of it and writes scale too. Most large designs use both, and the interview turns on how you choose the shard key and what you promise about reads that may be slightly behind.
Building block 5 of 6 in system design building blocks
When it comes up
- Data or write volume beyond one machine, according to your estimate.
- Read-heavy load that read replicas can absorb.
- Disaster recovery or lower latency in another region.
- Hot spots: one tenant, account or day getting far more traffic than the rest.
The core idea
Replication usually has one leader that takes writes and followers that copy them. Asynchronous copying keeps writes fast but leaves followers slightly behind, so a user can write and then read an older value from a replica; route that user's next reads to the leader, or read from a replica only once it has caught up. If the leader fails, a follower is promoted, and with asynchronous replication the last few writes can be lost; say whether that is acceptable.
Sharding picks a key and maps each key to a shard. Hashing the key spreads load evenly but makes range queries touch every shard; ranges keep nearby keys together but invite hot spots; sharding by tenant keeps each tenant's queries on one machine. Consistent hashing places shards on a ring so adding one moves only a small share of the keys. Queries and transactions that span shards are slow and complex, so the best key is the one most queries already include.
A sketch: consistent hashing
import bisect, hashlib
def _point(key):
return int(hashlib.md5(key.encode()).hexdigest(), 16)
class Ring:
"""Each shard appears at many points, so load stays even and adding a
shard moves only the keys that now land on its points."""
def __init__(self, shards, points_per_shard=100):
self.ring = sorted((_point(f"{s}#{i}"), s)
for s in shards for i in range(points_per_shard))
self.positions = [p for p, _ in self.ring]
def shard_for(self, key):
i = bisect.bisect(self.positions, _point(key)) % len(self.ring)
return self.ring[i][1] # the first shard clockwise from the keyTrade-offs to name
- Hash keys (even load) against range keys (efficient range queries, hot spots).
- Asynchronous replication (fast writes, possible loss on failover, stale reads) against synchronous (slower writes, no loss).
- More shards (headroom) against operational cost and more cross-shard work.
- Sharding now against later: if one machine plus replicas will do for years, say so and plan the migration instead.
Common mistakes
- A shard key that creates hot spots, such as today's date or one very large customer.
- No plan for adding shards or moving data while the system stays up.
- Reading from a replica straight after a write and showing the user stale data.
- No failover plan, or one that silently loses acknowledged writes.
- Sharding because it sounds impressive when the estimate says one database is enough.
How to explain it out loud
Justify the shard key with the queries: "I'll shard by workspace, because nearly every query is scoped to one workspace, which keeps it on one shard." Then name the exception that breaks the rule, such as a workspace that outgrows a shard, and how you would handle it.
Say plainly what can be stale and for how long, and what happens on failover. Scalability is a quarter of the design score in Devana's rubric and reliability a fifth; this topic touches both, and candidates who state the consistency they are giving up score better than those who leave it implied.
Practice questions
These come from Devana's question bank, in the order to try them. Each one starts a voice mock interview with Josh, Devana's AI interviewer, on that question, so you practice explaining the approach out loud as well as getting it right.
- Practice
Design sharding for workspace data
mediumNotion · Choosing a key for locality, and moving data safely.
- Practice
Design partitioning for a very large table
mediumOracle · Range partitions aligned with the main filter.
- Practice
Design cross-region replication for a database
mediumSnowflake · Lag, failover and failback across regions.
- Practice
Design replication for a database
hardMongoDB · Leaders, followers and what an acknowledged write means.