Mid to Senior Engineer

System Design Interview Prep

A structured path from the interview framework through core concepts, key technologies and patterns to eighteen full problem breakdowns, each with diagrams and weak, solid and excellent answers to every deep dive.

Chapter 12 of 36Patterns · Scaling Reads

Pattern: Scaling Reads

Most products are read-heavy. A social app might serve a hundred feed views for every post, and an online shop serves thousands of product views for every order. When a design asks "how do you handle ten times the traffic?", the answer for the read side is nearly always a combination of the techniques in this chapter. Learn them as a toolbox, with the order in which you reach for each, and the cost each one imposes.

The unifying idea: do less work per read, or do the work once and reuse it. Every technique either moves reads to something cheaper, or removes them entirely.

1. Start by measuring where the reads go

Before choosing a tool, know the shape of the read load:

  • The read to write ratio. A ratio of 100 to 1 invites heavy caching and replicas. A ratio of 2 to 1 makes both less useful, because writes invalidate and replicate constantly.
  • Skew. If 20 percent of items get 80 percent of reads, a small cache catches most traffic. If reads are spread evenly across a huge dataset, caching helps little and you need capacity instead.
  • Freshness requirement. How stale may a read be: a second, a minute, never? Staleness tolerance is what you spend to buy speed.
  • Query shape. Lookups by key are cheap to scale. Complex queries, joins and aggregations are expensive, and need different treatment.

A strong opening is: "I would first check the read to write ratio, how skewed the access is, and how fresh reads must be, because those decide which of the following I need."

2. The layers, from cheapest to most expensive

Think of read traffic as passing through layers, each absorbing a share.

<!--fig:reads-->
misses misses fresh only All reads 100% Browser + CDN absorbs most static App cache hot objects Read replicas scale-out reads Primary writes + the rest Also: denormalised read models, precomputed results, search andanalytics replicas. Each trades freshness for speed. Figure 1. Layers that shrink read traffic. Percentages are illustrative: each layer absorbs most of what reaches it.

Layer 0: do not make the request

  • Client-side caching. Browsers and apps cache responses according to Cache-Control and entity tag headers. A revalidation request that returns "not modified" costs far less than a full response.
  • Request coalescing and deduplication in the client. Do not fetch the same thing twice on one screen.
  • Batching. One request for many items beats many requests.
  • Precomputation in the client of anything that does not need the server.

Layer 1: CDN and edge

For content shared by many users, such as images, scripts, public pages and cacheable API responses, serve it from an edge location near the user. This removes the traffic from your servers entirely, and cuts latency. Use versioned URLs so you can cache for a long time and still publish changes, and use short lifetimes with revalidation for content that changes. Do not cache personalised responses at a shared edge.

Layer 2: application cache

A shared in-memory cache, as in the Redis chapter, holds hot objects and expensive query results. Use cache-aside with a time to live. Decide the invalidation rule up front: expire by time, or delete on write. Protect against stampedes and hot keys, as described in the scaling and distributed cache chapters.

Cache what is expensive and shared. A cache hit on a trivial primary-key lookup may save little. A cache entry for an aggregate that takes 300 milliseconds to compute saves a lot, and serves many users.

Layer 3: read replicas

Copy the database to read-only replicas and send queries to them. You scale read throughput by adding replicas, and you also gain availability.

The cost is replication lag. A replica may be seconds behind. Techniques to handle it:

  • Read your own writes. After a user writes, read from the primary for a short window or until the replica reaches that position.
  • Sticky routing. Send one user's reads to one replica, so they never see time go backwards.
  • Choose by tolerance. Send reads that tolerate staleness (a product page, a feed) to replicas, and reads that must be current (an account balance before a payment) to the primary.

Layer 4: the primary

What is left. It should now carry the writes and only the reads that must be fresh.

3. Reshape the data for reads

When queries are expensive, no amount of caching fixes the underlying cost. Change how the data is stored.

Indexes. The cheapest fix. A missing index turns a lookup into a scan. Check query plans before adding any infrastructure.

Denormalisation. Store the data in the shape the read needs, duplicating it if necessary. A product listing page that needs the product, its price, its brand name and its average rating should not join four tables per request. Store a precomputed row or document with those fields, and update it when the source changes. The cost is write complexity and the risk of copies disagreeing, so decide which copy is authoritative and how the others are refreshed.

Precomputation and materialised views. Compute expensive aggregates ahead of time: daily totals, top-N lists, per-user feed timelines. Reading becomes a lookup. Refresh on a schedule or incrementally on change. This is the idea behind the news feed's precomputed timelines and the typeahead's precomputed suggestions.

Specialised read stores. Feed a search engine for text queries, a column-oriented store for analytics, or a graph store for relationship queries, from the primary database through change capture. Each is a read model optimised for one kind of question, and each is eventually consistent with the source.

CQRS in one line. Separate the model you write to from the models you read from. It is the general name for the ideas above, and you pay for it with eventual consistency and more moving parts.

4. Scale the data layer horizontally

If reads are spread across a dataset too large or too busy for one machine's replicas, shard it. Choose a key so that the dominant read touches one shard, as in the storage chapter. A wide-column or partitioned key-value store is built for this: a lookup by partition key goes to one node whatever the cluster size.

Sharding for reads is the heaviest option. Reach for it only after caching, replicas and data reshaping, and when the numbers show that they are not enough.

5. Special cases

Hot keys and viral content. One item gets enormous traffic, more than a single cache node or shard can serve. Replicate the item across several cache nodes, add a short-lived local cache in each application server, or serve it from the edge. Detect it from access sampling.

Cold start and cache warm-up. A new cache or a new region starts empty, and the first wave of requests all miss. Pre-populate the cache before sending it traffic, or ramp traffic up gradually, so the database is not flattened.

Read-heavy and very large objects. Put large, immutable objects in object storage with a CDN in front. Keep only a reference in the database.

Geographic reads. Place replicas and caches in the regions where users are, so reads never cross an ocean. Accept that writes may still go to one region, and that cross-region replication adds lag.

Rate limiting and load shedding. When reads exceed even a scaled system, protect the core by limiting per client and shedding low-priority requests, so that important reads still succeed.

6. Choosing among them: a worked example

Problem. A product catalogue page is slow and the database is at 90 percent CPU. Traffic is 20,000 reads per second, 200 writes per second, and 80 percent of reads hit 5 percent of products.

Reasoning.

  1. The read to write ratio is 100 to 1, and access is heavily skewed. That is the profile for caching.
  2. Check the query: the page joins five tables. Add the missing indexes and denormalise the page's data into one precomputed record per product.
  3. Put the product page behind a CDN with a short lifetime, and cache the record in an application cache with a time to live of a minute. A 90 percent hit rate cuts database reads from 20,000 per second to about 2,000.
  4. Add one or two read replicas for the remaining reads and for search, accepting a second or two of lag.
  5. A price change should appear quickly. Delete the cache entry on write, and purge the edge for that product, with the short TTL as the safety net.

What I would say about the trade-off. "A customer might see an old price for up to a minute at the edge. For this catalogue that is acceptable, and at checkout I read the price from the primary, so the charge is always current."

7. Interview questions and model answers

Q: A read-heavy service is overloaded. What do you do first? Measure: the read to write ratio, the skew of access and the freshness required. Then fix the cheapest things first: indexes and query shape, then caching, then replicas, then data reshaping, and sharding last.

Q: How do you handle replication lag? Read from the primary for data the user just wrote or that must be current, use sticky routing so a user does not see time go backwards, and route staleness-tolerant reads to replicas.

Q: When does caching not help? When access is spread evenly over a huge dataset, so the hit rate is low, or when the data changes so fast that entries are invalidated before they are reused, or when the freshness requirement forbids staleness.

Q: How do you serve a viral item? Replicate the hot key across cache nodes, add a small local cache in each server and serve it from the CDN, so no single node carries it.

Q: What is denormalisation and what does it cost? Storing data in the shape a read needs, duplicating fields to avoid joins. It makes reads fast and writes more complex, and copies can disagree, so I define the source of truth and how the others are refreshed.

8. Common mistakes

  • Adding infrastructure before checking indexes and query plans.
  • Caching without a clear invalidation rule.
  • Sending every read to replicas, including ones that need current data.
  • Ignoring skew, or assuming a cache will help on evenly spread data.
  • Forgetting cache warm-up and cold-start load on the database.
  • Sharding for reads before exhausting caching, replicas and reshaping.
  • Caching personalised data at a shared edge.
Header Logo