Data Storage: Choosing a Database, Replication and Sharding
Almost every system design question ends up as a question about where the data lives and how it is split. This chapter covers how to choose a store, how to keep copies of it, how to divide it across machines and how to reason about what you give up.
1. Choosing a store from the access pattern
Do not start from the brand. Start from how the data is read and written.
| Store type | Strength | Typical use | Watch out for |
|---|---|---|---|
| Relational (SQL) | Transactions, joins, flexible queries, strong schema | Orders, accounts, anything with relationships | Scaling writes beyond one primary needs work |
| Key-value | Very fast lookup by key, easy to partition | Sessions, caches, counters | No rich queries |
| Document | Nested records, flexible schema | Catalogues, profiles | Joins are awkward, schema drift |
| Wide-column | Huge write throughput, time-ordered data, partitioned by design | Event logs, time series, messaging | Queries must match the table's design |
| Graph | Traversals over relationships | Social graphs, recommendations | Not for bulk analytics |
| Search index | Full-text and faceted search | Product search, logs | Not the source of truth |
| Object storage | Cheap, durable blobs | Images, video, backups | Not for small random updates |
A good habit: name the dominant query. "Fetch a record by its short code" points to a key-value or primary-key lookup. "Show me the last 50 messages in this conversation" points to a store partitioned by conversation and ordered by time.
2. SQL versus NoSQL, honestly
The usual slogans are misleading. SQL databases scale far further than people assume, and many NoSQL stores now offer transactions. Use these questions instead:
- Do I need multi-row transactions and joins? Lean relational.
- Is the access pattern simple and known, with enormous volume? A key-value or wide-column store may suit.
- Will the schema change constantly? A document model may reduce friction.
- Do I need to query by many different attributes? That favours relational or a secondary search index.
A common, defensible answer is polyglot: a relational database as the source of truth, a cache in front of it, and a search index fed from it.
3. Indexes
An index is a separate structure that lets the database find rows without scanning the table. The usual implementation is a B-tree, which keeps keys sorted, supports equality and range queries, and costs per lookup.
Trade-offs to state:
- Each index speeds reads and slows writes, because every insert and update must maintain it.
- Indexes use memory and disk.
- A composite index on
(a, b)serves queries that filter ona, or onaandb, but not onbalone. Column order matters. - An index on a low-cardinality column, such as a boolean, rarely helps.
Write-heavy stores often use a log-structured merge tree (LSM tree) instead: writes go to an in-memory table and are flushed as sorted files that are merged in the background. Writes become sequential and fast, while reads may check several files, which is why they use Bloom filters to skip files that cannot contain the key.
<!--fig:lsm-->4. Replication
Replication keeps copies of the data on several machines for availability and read capacity.
Leader-follower (primary-replica). All writes go to the leader, which streams changes to followers. Reads can be served by followers. This is the default and the simplest to reason about.
- Synchronous replication waits for the follower to confirm. No acknowledged write is lost on failover, but latency rises and a slow follower stalls writes.
- Asynchronous replication acknowledges after the leader writes. It is fast, but a leader crash can lose the last few writes, and followers may serve stale data.
- A common compromise is semi-synchronous: wait for one follower, let the rest lag.
Replication lag causes visible anomalies. A user posts a comment, reloads, and the follower has not received it yet. Fixes: read your own writes from the leader for a short window, route a user's reads consistently to one replica, or attach a version to the session so a replica waits until it has caught up.
Failover. When the leader dies, a follower is promoted. Hard parts: detecting failure without false alarms, choosing the most up-to-date follower, and preventing the old leader from continuing to accept writes after it returns (a split brain). Consensus-based systems solve this by requiring a majority to agree on who the leader is.
Multi-leader and leaderless designs accept writes in several places. They improve write availability and cross-region latency, and they introduce write conflicts that must be resolved by timestamps, application logic or conflict-free data types.
5. Sharding (partitioning)
When one machine cannot hold or serve the data, split it across several. Each piece is a shard.
Choosing the shard key
The key decides everything:
- It must spread data and load evenly.
- It should keep related data together so common queries touch one shard.
- It should be present in the common queries.
A bad key creates a hot shard: sharding by creation date sends all new writes to one shard, and sharding a social network by celebrity concentrates traffic. A user identifier is often a good key for per-user data, though a very large user can still be hot.
Strategies
- Range partitioning splits by key ranges. It supports range queries, but sequential keys create hot spots.
- Hash partitioning applies a hash to the key and assigns it to a shard. It spreads load evenly but loses range ordering.
- Directory-based partitioning keeps a lookup table from key to shard. It is flexible, and the table becomes a component you must make reliable.
Consistent hashing
With plain hash(key) mod N, changing remaps almost every key, forcing a massive data move. Consistent hashing places both nodes and keys on a ring. A key belongs to the first node clockwise from its position. When a node is added or removed, only the keys between it and its neighbour move, which is about of the data on average.
Real systems add virtual nodes: each physical node owns many points on the ring. This evens out the distribution and lets a stronger machine take more points.
The costs of sharding
- Cross-shard queries and joins are expensive or impossible, so the design avoids them.
- Cross-shard transactions need coordination such as two-phase commit, or are replaced by sagas with compensating actions.
- Resharding is operationally hard, so plan capacity or use a scheme that splits shards gradually.
- Secondary indexes across shards need their own design.
A sensible interview line: "I would avoid sharding until the numbers force it, because it makes everything harder. When I do shard, I would pick the key from the dominant query."
6. The CAP theorem and what it really says
The CAP theorem states that during a network partition, a distributed store must choose between consistency (every read sees the latest write or an error) and availability (every request gets a non-error response). Partitions happen, so the real choice is made when one occurs.
Important corrections that mark you out:
- It is not "pick any two of three" as a permanent label. When there is no partition, you can have both consistency and availability.
- "Consistency" in CAP means linearizability, which is stronger than the "C" in ACID.
The PACELC refinement adds the normal case: if there is a Partition, choose A or C; Else, choose Latency or Consistency. Even without failures, waiting for replicas to agree costs latency.
7. Consistency levels you will be asked about
- Strong (linearizable): reads see the latest committed write.
- Eventual: replicas converge if writes stop, but reads may be stale meanwhile.
- Read-your-writes: a client sees its own updates.
- Monotonic reads: a client never sees time go backwards.
- Causal: effects are never seen before their causes.
Say which one the feature needs. A like count can be eventual. A bank transfer cannot.
Potential deep dives
Deep dive 1: How do you choose the database?
The challenge. The prompt describes a product, and you must commit to a store and defend it.
Weak: choose by popularity or habit. "I will use a NoSQL database because it scales." That is a slogan, not an argument, and it hides that most NoSQL stores trade away joins and transactions the product may need.
Solid: choose from the access pattern. List the main queries, the read to write ratio, the consistency needs and the data shape. Relational when you need transactions, relationships and flexible queries. A key-value or wide-column store for very high write volume with known lookups. A search index for text. Object storage for blobs.
Excellent: choose, justify and plan the evolution. Start from the simplest store that meets the numbers, name the point at which it stops being enough, and say what you would do then: "A relational database with a read replica and a cache handles this load for years. If writes grow past one primary, I would partition by tenant, and if search is needed, I would feed a search index from the database through change capture." Mention polyglot persistence honestly, one source of truth and derived read models, and the cost of keeping them in sync.
Deep dive 2: How do you choose a shard key and handle resharding?
The challenge. One machine cannot hold or serve the data, so you must split it, and the split is hard to change later.
Weak: shard by an auto-incrementing identifier, with modulo. Sequential ids send all new writes to one range if range partitioned, and modulo hashing remaps almost everything when the shard count changes.
Solid: hash a key that appears in the main query, with consistent hashing. A hash spreads load evenly. Consistent hashing, or a fixed set of virtual partitions assigned to nodes, moves only a small share of data when you add a node. Keep related data on one shard so the main queries touch one shard.
Excellent: plan for growth, skew and online migration. Create many more logical partitions than nodes at the start, so that growth means moving whole partitions between nodes, not splitting data. Monitor per-partition load and size and split or move hot ones. For a live migration: copy the partition in the background, stream changes to the new copy, verify it, switch routing atomically through a shard map, and keep the old copy briefly for rollback. Address skew directly: a very large tenant or celebrity key can need its own shard or a salted key. Say that cross-shard queries and transactions are the price, and that you design the data model to avoid them in the main flows.
Deep dive 3: How do you handle primary failure without losing data?
The challenge. The primary dies. Writes must continue, and acknowledged writes should survive.
Weak: promote any replica. An asynchronous replica may be seconds behind, so promoting it silently loses acknowledged writes, and if the old primary returns and keeps accepting writes you have two diverging copies.
Solid: promote the most up-to-date replica, and fence the old primary. Choose the replica that has applied the most of the log, update routing so clients find it, and ensure the old primary cannot accept writes if it returns, for example by requiring it to rejoin as a replica.
Excellent: decide the durability trade-off per data class. For data that must not be lost, use synchronous or quorum replication, so that a failover never loses an acknowledged write, and pay the latency. For data that tolerates a small loss, use asynchronous replication. Use a majority-based election to prevent split brain. Test failover regularly, measure recovery time and data loss in each scenario (recovery time objective and recovery point objective), and keep backups and point-in-time recovery for the cases replication cannot fix, such as a bad deployment that corrupts data on every replica.
Deep dive 4: How do you handle replication lag?
The challenge. Users read from replicas that lag behind the primary, and see their own changes disappear.
Weak: ignore it. Users post a comment, reload, and it is missing. They post it again.
Solid: read your own writes from the primary. After a write, route that user's reads to the primary for a short window.
Excellent: consistency tokens and targeted routing. Return a position token (a log sequence number or timestamp) with each write. The client sends it with later reads, and a replica serves the read only if it has reached that position, otherwise the request waits briefly or goes to the primary. Add sticky routing so a user's reads come from one replica, preventing time from going backwards. Reserve primary reads for the cases that need them, so the primary is not overloaded.
What is expected at each level
Mid-level. You pick a reasonable store for the data, add replication for availability and read capacity, and explain what sharding is and why it is needed.
Senior. You choose a store from access patterns, a shard key from the dominant query, explain consistent hashing and its virtual nodes, handle failover and replication lag, and state the CAP and consistency trade-offs for the data.
Staff. You discuss the evolution of the data tier over years, online resharding, per-data-class durability, backup and recovery objectives, cost of storage tiers, and the operational burden of each choice.
Interview questions and model answers
Q: How would you scale a relational database that is read-heavy? Add read replicas and route reads to them, accepting replication lag, and put a cache in front for the hottest queries. Check indexes and slow queries before adding hardware.
Q: And when it becomes write-heavy? First optimise: batch writes, remove unneeded indexes, move large blobs to object storage, and queue non-urgent writes. When one primary is truly saturated, shard by a key chosen from the dominant query, and avoid cross-shard transactions by design.
Q: Why use consistent hashing? So that adding or removing a node moves only a small fraction of keys, roughly one over the number of nodes, rather than reshuffling everything. Virtual nodes keep the load even.
Q: Your primary fails. What happens to writes in flight? It depends on replication mode. With asynchronous replication the latest acknowledged writes may be lost. With synchronous or quorum replication they are safe but each write is slower. I choose by how costly lost data is for this feature.
Q: How do you pick a shard key for a chat application? Conversation identifier. Messages for one conversation stay on one shard, ordered by time, so the main query touches one partition. I handle a very large group chat as a hot-shard case, for example by splitting it into time buckets.
Common mistakes
- Choosing a database by popularity.
- Sharding early, with a key that is absent from the main queries.
- Forgetting replication lag when serving reads from followers.
- Quoting CAP as "choose two" with no mention of partitions.
- Using hash modulo N and ignoring resharding.