System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Scaling the data: replication, caching, CDN and sharding


The web tier now scales by adding machines. The database still does not: there is one, and it is the single point of failure for the whole system.

This lesson continues growing Trailmix, the route-sharing app for cyclists, from where the stateless web tier left it. Four moves follow, each forced by the one before. Replication scales reads and removes the database as a single point of failure. A cache takes most reads off the database entirely. A CDN moves static bytes close to users. And sharding, the most expensive step, is the only one that scales writes.

Leader and followers

Replication means keeping copies of the same data on several machines. The standard arrangement is one leader (also called primary or master) and several followers (replicas or read replicas).

  • All writes go to the leader.
  • The leader records every change in an ordered log and ships it to the followers.
  • Each follower applies the changes in the same order, so it converges to the same state.
  • Reads can go to any follower.
APPLICATIONApp serversPRIMARYLeaderall writesREPLICASFollower 1readsFollower 2readswritesreadsreadsreplication log1. how stale may a read be? replication lagis usually milliseconds and occasionallyseconds2. what happens when the leader dies?promotion is a decision, not an automaticfact
Replication buys read capacity and costs consistency — the two annotations are the questions an interviewer will always ask next.

What replication buys

Read scaling. Trailmix is read-heavy: browsing routes vastly outnumbers posting them. Say 95% of queries are reads. One leader plus three followers moves 95% of the query load off the leader, so the leader's load drops roughly twenty-fold. Need more read capacity? Add a follower. This is the cheapest scaling move in this section.

Availability. If a follower dies, the balancer stops sending it reads. If the leader dies, a follower can be promoted.

Backups and analytics. Run the expensive nightly report against a follower so it does not compete with user traffic.

Replication lag — the first thing every interviewer probes

Replication is usually asynchronous: the leader acknowledges a write as soon as it is durable locally, without waiting for followers. Lag is typically single-digit milliseconds within one datacentre, but it grows to seconds under heavy write load, during a long transaction, or while a follower is catching up after a restart. Across continents, the network round trip alone is 100–200 ms.

The user-visible symptom is precise: a cyclist uploads a route, the write goes to the leader, the app immediately reloads the profile page, that read goes to a follower that has not caught up, and the route they created 300 ms ago is not there. They upload it again.

Three standard fixes:

  1. Read your own writes from the leader. After a user writes, route that user's reads to the leader for a short window (say 10 seconds). Simple, effective, and the usual answer.
  2. Read from the leader for data the user can modify, and from followers for everything else. Coarser, cheaper to implement.
  3. Track the write position. The client keeps the log position of its last write and the follower serves the read only when it has reached that position. Precise and fiddly.

Consistency models, in Section 5, names this guarantee: read-your-writes consistency.

When the leader dies

This is the harder half, and it is where candidates run out of material.

A follower must be promoted. Something has to decide which — usually the most caught-up one — and every other node and client must learn about the change. During the switch, writes fail; a fast automatic failover is on the order of 10–30 seconds.

Two problems worth naming unprompted:

Lost writes. With asynchronous replication, writes acknowledged by the old leader but not yet shipped are gone when a follower is promoted. Semi-synchronous replication — the leader waits for at least one follower before acknowledging — closes this at the cost of adding that follower's round trip to every write.

Split brain. If the old leader recovers and does not know it was replaced, two nodes accept writes and the data diverges. Prevented by fencing: a majority-based coordinator issues a monotonically increasing term number, and the old leader's writes are rejected.

Caching

Replication adds read capacity by adding databases. A cache removes most of the reads instead. A cache is a small, fast store holding copies of data that is expensive to fetch. It is the highest-leverage component in this section and the one that introduces the most subtle bugs.

The cache-aside read pathRequest arrivesLook inthe cacheMiss: querythe databaseWrite backwith a TTLReturn tothe callerA TTL too short thrashes; too long serves stale data.
Cache-aside puts the miss path in your own code, which is exactly where the three cache bugs live.

Trailmix's most-requested endpoint returns a popular route with its statistics: four joins, about 5 ms on the database. Put a cache — an in-memory key-value store such as Redis or Memcached — between the application and the database, and use the cache-aside pattern:

  1. Application asks the cache for key route:8213.
  2. Hit: return it. Roughly 0.5 ms including the network round trip.
  3. Miss: query the database (5 ms), write the result into the cache with a time-to-live (TTL), return it.

The arithmetic that justifies it

At a 95% hit ratio, average latency is (0.95 × 0.5) + (0.05 × 5) = 0.72 ms, against 5 ms before — a seven-fold improvement.

The bigger win is the database. At Trailmix's 1-million-user peak of 1,700 requests per second, the database was serving 1,700 queries per second. With a 95% hit ratio it serves 85. That is the difference between needing eight read replicas and needing one.

What to cache, and for how long

Cache things that are read far more often than they change, and where slightly stale is acceptable. Route details, user profiles, follower counts: yes. A payment balance: no (the digital wallet in Section 29 explains why).

TTL is a direct trade of freshness against load. A 60-second TTL means data can be up to 60 seconds stale and the database sees at most one query per key per minute. Too short and the hit ratio collapses; too long and users see stale data for longer than they will tolerate.

Eviction handles a full cache. Least Recently Used (LRU) — discard the item untouched for longest — is the sensible default and matches how access is usually distributed. Caching strategies, in Section 5, covers the alternatives.

The three problems that come free with every cache

1. Staleness. A route is edited, the database is updated, and the cache still holds the old copy until its TTL expires. Fixes: delete the key on write (invalidation), write to both (write-through), or accept the TTL window. Deleting on write has its own race — a concurrent read can repopulate the old value between the database write and the delete.

2. The thundering herd on a hot key. A route goes viral and is being read 5,000 times a second. Its cache entry expires. Every one of those 5,000 requests misses within the same few milliseconds and hits the database, which was sized for 85 queries per second. The database saturates, requests time out, and the cache never gets repopulated because nothing completes. The fix is request coalescing: the first miss takes a short lock and fetches; the others wait for its result. Serving the expired value while one request refreshes in the background also works well.

3. The cold start stampede. You restart the cache cluster. Every key is missing at once. The database receives the full 1,700 requests per second it has not seen since the cache was added, plus the backlog, and falls over. Mitigations: warm the cache before taking traffic, jitter TTLs so keys do not expire in lockstep, and rate-limit the fill path.

CDN and static content

The cache fixes database load for data. Photos, scripts and other static files have a different problem: distance. A Content Delivery Network (CDN) is a fleet of servers spread across the world that hold copies of your static files and serve them from wherever the user is.

How far the bytes have to travelUser's browserEdge PoP: 10 ms awayRegional cacheOrigin: 150 ms away
Static files are most of the bytes and none of the logic, so distance is the only thing left to optimise.

Trailmix stores its route photos in one region — say Mumbai. A user in Berlin requests a photo. The round trip is roughly 130–150 ms of pure network latency, before the file starts transferring. For a 500 KB photo over a mediocre connection, the total is comfortably over a second, and a route page loads several.

With a CDN, the first Berlin request fetches the photo from your origin and stores it at the Frankfurt point of presence. Every later Berlin request is served from about 10–20 ms away. The photo also stops consuming your origin's bandwidth.

The volume that makes a CDN non-negotiable

At 10 million daily active users, 2% of whom post a route with three photos a day, that is 200,000 routes and 600,000 photos daily. At 500 KB each, 300 GB of new photos a day, or about 110 TB a year.

Serving is the larger number. If each photo is viewed 50 times, that is 30 million photo requests a day, or 15 TB a day of egress — around 1.4 Gbps sustained, several times that at peak. Pushing that through your own servers means provisioning for peak bandwidth you use for four hours a day.

What belongs on a CDN

  • Images, video segments, fonts, JavaScript, and stylesheets. Always.
  • Generated files that are the same for everyone: an export, a public route map tile.
  • Cacheable API responses. A public route's JSON is identical for every viewer and can sit at the edge for 60 seconds. Anything personalised cannot, unless you vary the cache key.

Never put user-specific responses behind a shared cache without a key that includes the user, and treat that as a security question rather than a performance one. Serving one user's private data to another from an edge cache is a well-known class of incident.

CDN invalidation, done properly

You have published logo.png with a one-year TTL, and now the logo has changed. Two options.

Purge. Ask the CDN to drop the object everywhere. It works, it takes seconds to minutes to propagate, and at high frequency it is slow and sometimes rate-limited.

Versioned URLs. Put a hash of the content in the filename — logo.9f3ac1.png — and give it an immutable, one-year TTL. Changing the logo produces a different filename, referenced by an HTML file with a short TTL. Nothing is ever invalidated because nothing is ever overwritten.

Recommendation: versioned URLs for everything you deploy; purge reserved for the case you cannot version, such as removing content that must disappear immediately.

The CDN cost model

CDN egress is usually cheaper per gigabyte than cloud origin egress, and the saving is secondary. The real saving is capacity: you no longer size your origin for peak global traffic. What you pay in exchange is a cache-hit-ratio problem — a CDN with a poor hit ratio costs you both CDN charges and origin bandwidth. Watch that ratio; the usual culprits are unnecessary query strings in URLs and cache headers that prevent storage.

Sharding: when it becomes necessary

Replication scales reads, the cache absorbs most of them, and the CDN takes the static bytes away. Nothing so far scales writes, because every write still goes to one leader. Sharding is how you fix that, and it is the most expensive step in this section.

Trailmix at 10 million daily active users: 200,000 routes a day, each about 9 KB of track points and metadata, is only 1.8 GB a day and 2.3 writes per second. That does not need sharding.

The activity table does. Every view, like, and comment is recorded: say 10 million users generating 30 events a day, so 300 million rows a day at 200 bytes — 60 GB a day, 22 TB a year, and 3,500 writes per second sustained. One leader can absorb a few thousand small writes per second, so you are at the edge, and the dataset outgrows a single machine's disk within a year.

What sharding is

Split one logical table across several databases (shards), each holding a disjoint subset of rows and each with its own leader and followers. A shard key decides which shard a row lives on.

Two ways to map keys to shards, plus a third worth knowing:

  • Hash-based. shard = hash(user_id) % 16. Even distribution, no range scans. The modulo form is a trap when the shard count changes — Section 7 (Design Consistent Hashing) is entirely about that.
  • Range-based. Shard by route_id ranges or by date. Range scans stay efficient, but sequential keys create a hot shard: if you shard by date, today's shard takes every write.
  • Directory-based. A lookup service maps each key to a shard. Maximum flexibility, at the cost of one extra hop and a component that must not go down.

For Trailmix's activity table, hash on user_id — activity is nearly always queried per user, which keeps a user's rows on one shard.

STATELESS TIER — scales horizontallyCLIENTSWebMobileEDGECDNLoad balancerAPI serversN identical instancesCacheread-throughMessage queueAsync workersemail, thumbnailsDATAPrimary DBwritesReplicasreadsObject storestaticAPIround robinmisswritereadstorereplicateorigin pullnothing here holds session state, so anyinstance can serve any requestthe only stateful layer, and the onlyone that is hard to scale
Read the layers left to right: everything before the data tier is disposable, which is exactly why it scales.

What sharding takes away

This is the part that earns marks.

Cross-shard joins stop working. "Routes from people I follow" spans every shard the followed users live on. You either query all sixteen shards and merge in application code, or denormalise so the answer lives on one shard. Both are real work you did not have before.

Multi-row transactions stop working. A transaction spanning two shards needs distributed coordination. Two-phase commit is available, slow, and blocks on a coordinator failure. In practice, systems avoid it — the hotel reservation system in Section 24 uses the saga pattern instead, and the payment system in Section 28 explains why a distributed transaction is usually the wrong answer for money.

Global uniqueness stops working. Auto-incrementing primary keys collide across shards. This is precisely why Section 9, on unique ID generation, exists.

Aggregates get expensive. COUNT(*) on one table becomes sixteen queries and a sum, and the answer is inconsistent because the shards are read at slightly different times.

Resharding is painful. Going from 16 shards to 32 with a modulo scheme moves roughly half of a multi-terabyte dataset while the system is live. Consistent hashing (Section 7) reduces the moved fraction to about 1/N, which is the main reason it exists.

Hot shards happen anyway. If one account is followed by ten million people, its shard gets a disproportionate share of traffic no matter how good the hash is. The fix is always specific to the workload — The hybrid that real systems use, in Section 13, handles exactly this case for feeds.

The order of moves, and why order is the lesson

Single server → split the database → load balancer and stateless web tier → replication → cache → CDN → shard. Each step was forced by a specific failure of the step before it, and each cost something. That sequence, with the reason for each move, is a better answer to "how would you scale this?" than any finished diagram.