Practice this topic in a realistic system design interview
In many system design problems, reads happen far more often than writes. For example: A short link might be created once but clicked millions of times. A social post might be written once and then read by thousands or even millions of followers.
This lesson covers how to design a system that handles a very large number of reads, the main techniques used to reduce load on the database, and the trade-offs that come with each one.
It also covers how to approach high read traffic in a system design interview, since many common interview problems are read heavy.
Let's start with the kinds of problems where reads are much more common than writes.
| Problem | Written | Read |
|---|---|---|
| URL shortener (Bitly) | A short link is created once | It can be clicked millions of times |
| Design Twitter or Instagram | A post is written once | It's read by every follower who opens their feed |
| Design YouTube | A video's title, description, and view count change rarely | They're loaded every time someone opens the video page |
| E-commerce (Amazon) | Product details are updated occasionally | Product pages are viewed far more often |
| Ticket booking (Ticketmaster) | A ticket is bought once | Thousands of users keep checking the event page long before they buy |
In many of these systems, the ratio of reads to writes is 100:1 or even higher. For popular content, it can be much higher.
So the main challenge is serving all those reads quickly, without overloading the database.
Suppose we're designing a URL shortener that handles 1 billion redirects per day.
That's about 11,500 reads per second on average. During peak hours, traffic could be three times higher, around 35,000 reads per second.
Now imagine every one of those requests goes straight to a single database. Each query uses CPU, memory, and sometimes disk reads on that database server.
Depending on the hardware and the query, a single database server can usually handle somewhere between a few thousand and a few tens of thousands of simple queries per second. So at this scale, one database is not enough. Queries start to slow down, response times go up, and eventually requests start timing out.
Throughput is not the only concern. User-facing reads also have a latency budget: a redirect or an API call is expected in tens of milliseconds, and a database query that waits in a queue behind thousands of others will miss it. Cost matters too. Database capacity is expensive because it bundles CPU, memory, storage, replication, and backups, while serving a repeated read from memory or a CDN costs much less per request.
The goal of every technique in this lesson is the same: serve most reads without sending them to the main database.
Before adding new components, it's worth checking whether the database is being used efficiently.
The first thing to check is indexes.
If the URL shortener looks up links by short code, and there's no index on that column, every lookup has to scan the whole table. With an index, the same lookup takes only a few reads, even when the table has billions of rows.
For queries that filter on several columns, a composite index works the same way. Put the filter columns first, then the sort columns, so the index matches the query pattern.
The second thing to check is the queries themselves.
For example, a page that loads 20 posts and then runs a separate query for each post's author makes 21 queries instead of 2. This is called the N+1 query problem.
Fixing problems like this can reduce database load a lot, without adding any new infrastructure. EXPLAIN ANALYZE is the usual tool for finding them: it shows whether a query uses an index and how many rows it touches.
| Problem | Fix |
|---|---|
| Full table scan | Add an index on the lookup column |
| N+1 queries | Use a JOIN or batch fetch with IN (...) |
| Large result set | Add LIMIT and paginate |
| Expensive aggregation on every read | Precompute it (see section 10) |
In an interview, you usually don't need to spend much time here, but it's good to mention that the right indexes are in place before you move on.
Once the queries are efficient, another way to handle more reads is to spread them across more machines. This is where read replicas come in.
A read replica is a copy of the main database that stays in sync with it. All writes go to the main database, which is called the primary. The primary then sends those changes to one or more replicas. Reads can go to any of the replicas.
So if one database can handle 10,000 reads per second, adding four replicas gives you roughly 50,000 reads per second in total.
This works well for read-heavy systems because writes still go to only one place, while reads can be spread across many machines. The scaling is not perfectly linear, since every replica still has to apply every write, and each one adds storage and operational cost. Replicas also don't help when the load is on a single row. A popular post gets the same load on every replica, which is why a single popular item needs caching (section 8) rather than more replicas.
The routing can be done in the application, which sends writes and fresh reads to the primary and everything else to a replica, or in a database proxy such as ProxySQL that splits queries by type.
But replicas come with one important trade-off.
Copying changes from the primary to a replica takes time. Most of the time, it's a few milliseconds. But when the system is under heavy load, a replica can fall behind by seconds. This delay is called replication lag.
Lag exists because most systems replicate asynchronously: the primary confirms a write as soon as it's durable locally, without waiting for replicas. Synchronous replication removes most of the lag, but every write then waits on a replica, which raises write latency and lets a slow replica block writes. That's why asynchronous replication is the common default.
Let's look at an example. Suppose a user updates their profile picture.
The write goes to the primary. Then the page reloads, and the read goes to a replica that hasn't received the change yet. Now the user sees their old profile picture and thinks the update failed. This is called a read-your-own-writes problem.
One common fix is to send a user's reads to the primary for a short time after they make a change, for example, for the next few seconds.
Another option is to always read the user's own data from the primary, while reads of other users' data go to replicas. A third option is lag monitoring: track how far behind each replica is, and stop sending reads to it when the lag is too high for that endpoint.
| Fix | How It Works | Trade-off |
|---|---|---|
| Primary after a write | Route a user's reads to the primary for a few seconds after they write | Adds some load to the primary |
| Own data from primary | A user's own profile, settings, and orders always read from the primary | More primary load for data-heavy users |
| Lag monitoring | Remove replicas from rotation when they fall too far behind | Fewer replicas available under load |
For many features, like viewing someone else's post or a product page, a delay of a second or two is fine. So in an interview, it's worth asking which reads need to be up to date and which ones can be slightly stale.
Replicas add capacity, but every read still runs a full database query. Caching avoids that query for most reads, which is why it's usually the first layer to add.
A cache stores frequently read data in memory, usually in a system like Redis or Memcached. Reading from memory is much faster than querying a database, often under a millisecond.
The most common way to use a cache is called cache-aside. When a request comes in, the application first checks the cache. If the data is there, it returns it right away. If not, it reads from the database, stores the result in the cache, and then returns it.
The number that matters most here is the cache hit ratio, which is the percentage of reads served from the cache.
Let's go back to our URL shortener with 35,000 reads per second at peak. If the cache has a 95% hit ratio, only 5% of reads reach the database. That's about 1,750 reads per second, which a single database can handle easily.
Small changes in the hit ratio matter. At a 90% hit ratio the database sees 3,500 reads per second, twice as many. In an interview, stating the hit ratio you expect is more convincing than saying "add a cache".
This works especially well when a small set of data gets most of the traffic, like popular links, trending posts, or best-selling products.
Good candidates are read often, change rarely, and can tolerate a little staleness.
| Good Cache Candidates | Poor Cache Candidates |
|---|---|
| Short link to long URL mappings | Bank account balances |
| Product details | Prices that must be quoted exactly at checkout |
| Posts and profiles | One-time authentication tokens |
| Configuration data | Data that is rarely read |
Caching adds a new problem. Now the same data lives in two places, the database and the cache. When the data changes in the database, the cache may still have the old value.
There are two common ways to handle this.
The first is to set a time to live, or TTL, on each cache entry.
For example, if a product page is cached with a TTL of 5 minutes, the cached value expires after 5 minutes, and the next request loads a fresh copy from the database. This is simple, but users might see old data for up to 5 minutes.
The second approach is to delete the cache entry whenever the data changes.
When a product's price is updated, the application writes the new price to the database and then deletes that product from the cache. The next read gets a cache miss and loads the new price.
Deleting is usually safer than writing the new value into the cache. If two updates race, a cache write can leave the older value behind, while a delete just forces the next read to load whatever the database holds.
In practice, systems usually use both. They delete the cache entry on writes, and they also set a TTL, so that if a delete is ever missed, the old value still expires eventually.
Two other strategies come up in larger systems:
| Strategy | Staleness | Write Cost | Complexity | Good Fit |
|---|---|---|---|---|
| TTL only | Up to the TTL | None | Low | Data where a few minutes of staleness is fine |
| Delete on write + TTL | Very short, bounded by TTL if a delete is missed | One extra delete | Low | Most applications |
| Write-through | Very short | Higher, cache in the write path | Medium | Small read models that must stay fresh |
| Event-based | Short, depends on the queue | Low for the writer | High | Many services caching the same data |
There's another problem that shows up when a very popular cache entry expires.
Suppose a trending video's page is cached, and 50,000 users are loading it every second. When that entry expires, all of those requests get a cache miss at the same moment. And all of them go to the database to load the same data.
This sudden spike is called a cache stampede, and it can overload the database. There are a few ways to prevent it.
Let only one request rebuild the entry. Use a lock so that only one request loads the data from the database. The other requests either wait a moment and check the cache again, or get the old value in the meantime.
Refresh popular entries in the background. A background job reloads popular entries before they expire, so they never actually go missing.
Add randomness to TTLs. If many entries are created at the same time with the same TTL, they'll all expire together. Adding a few seconds of random variation spreads those expirations out.
| Technique | How It Works | Trade-off |
|---|---|---|
| Lock on rebuild | One request reloads the entry; the rest wait or get the old value | Waiting requests add a little latency |
| Background refresh | Popular entries are reloaded before they expire | Extra work for entries that are no longer read |
| TTL jitter | Random variation spreads expirations over time | Does not help a single hot key on its own |
Even when the cache is working well, a single key can get too much traffic.
In a distributed cache, each key is stored on one cache server. So if a celebrity with 100 million followers posts something, every read for that post goes to the same cache server. That one server becomes a bottleneck, while the other cache servers are mostly idle. This is called a hot key.
Adding more cache servers doesn't help, because the traffic is on one key, not spread across many.
One fix is to store copies of the hot key on several cache servers, for example, by adding a suffix like post:123:copy1, copy2, and copy3, and picking one at random on each read. The suffixed keys hash to different servers, so the reads are split across them. On a write, the application deletes every copy.
Another fix is to keep very popular data in a small in-memory cache on each application server. That way, most reads for the hot key never even reach the cache cluster.
The trade-off is that these local copies can be slightly out of date, so they usually have a very short TTL, like a few seconds.
A related problem is requests for keys that don't exist, such as a deleted post or a made-up short code. Each of these requests misses the cache and reaches the database. Caching the "not found" result for a short time, called negative caching, prevents this.
So far, every request still reaches our servers. For content that's the same for every user, we can stop many requests before they get that far.
A CDN, or content delivery network, stores copies of content on servers around the world, close to users. It's commonly used for images, videos, JavaScript, and CSS files. But CDNs can also cache full pages or API responses, as long as the response is the same for everyone.
For example, a product page on a public e-commerce site, or a news article, can be cached at the CDN for a short time. When a user in India requests that page, a CDN server in India can return it directly, without the request ever reaching your data center.
This reduces load on your servers, and it also makes pages load faster for users who are far from your data center.
The browser can cache content too. With headers like Cache-Control, the server tells the browser how long it can reuse a response before asking again. The same header also controls the CDN.
| Content | Example Cache-Control | Meaning |
|---|---|---|
| Versioned images, CSS, JS | public, max-age=31536000, immutable | Cache for a year; a new version gets a new URL |
| Public product page or article | public, max-age=60, s-maxage=300 | Browser reuses it for 60 s, the CDN for 300 s |
| User-specific data | private, no-store | Do not cache at the CDN or in the browser |
s-maxage applies only to shared caches like a CDN, which lets the CDN hold a response longer than the browser. Adding stale-while-revalidate lets the CDN serve the old copy for a short time after it expires while it fetches a fresh one in the background, so users don't wait for the refresh.
The main risk is caching something personal. Anything that differs by user, like a cart or account page, has to be marked private or split into a public part that the CDN caches and a small user-specific part fetched separately.
The techniques so far make each read cheaper by storing its result somewhere faster. Another approach is to do the expensive work at write time, so reads become simple lookups.
A good example is the home feed in Twitter.
One way to build the feed is to wait until the user opens the app, then look up everyone they follow, fetch recent posts from each of them, and sort the results. That's a lot of work for every single read.
Instead, when someone posts, the system can add that post to a precomputed feed for each of their followers. Then, when a user opens the app, the system just reads their feed from a single list. This is called fanout on write.
For accounts with millions of followers, writing to every follower's feed is too slow, so many systems use a hybrid: normal accounts fan out on write, and posts from very large accounts are merged into the feed at read time.
The same idea applies to counters. Instead of counting all the likes on a post every time it's viewed, the system keeps a like count and updates it whenever someone likes the post.
Materialized views and denormalized tables apply the same idea inside the database: store the result of an expensive join or aggregation, and read that instead of recomputing it.
This makes writes more expensive and uses more storage, but in a read-heavy system, it's usually a good trade-off, because reads happen far more often than writes.
Let's now look at how all of these layers work together.
Each layer handles part of the traffic, so only a small fraction of reads ever reach the database. Here's an example with the URL shortener's 35,000 reads per second at peak:
The exact percentages depend on how cacheable the traffic is, but the database ends up seeing a small fraction of the original load.
But each layer is also another place where data can be out of date. So when you add a layer, you should also think about how long its data can be stale, and how it gets updated.
| Layer | How Stale It Can Be | How It Gets Updated |
|---|---|---|
| Browser | Up to max-age | Expires, or a new versioned URL |
| CDN | Up to s-maxage | Expires, or an explicit purge |
| App server memory | A few seconds | Short TTL |
| Distributed cache | Up to the TTL, usually much less | Delete on write, plus TTL |
| Read replica | Replication lag, usually milliseconds | Replication stream from the primary |
Here's a step-by-step approach you can use in an interview.
1. Clarify the requirements. Ask about the read-to-write ratio, how many reads per second the system needs to handle, and how up to date the data needs to be.
| Question | Why It Matters |
|---|---|
| What is the read-to-write ratio? | Tells you whether reads are the main bottleneck |
| How many reads per second at peak? | Tells you whether one database is enough |
| How fresh must the data be? | Decides TTLs, replica use, and CDN caching |
| Is the response the same for every user? | Shared responses can go on a CDN; personal ones cannot |
| Can any single item suddenly get very popular? | Tells you whether you need to handle stampedes and hot keys |
2. Estimate the peak read load, so you know whether one database is enough.
3. Add caching, and explain what you're caching, how you keep it up to date, and roughly what hit ratio you expect.
4. Add read replicas if the database still can't handle the reads that miss the cache, and explain how you'll deal with replication lag.
5. Add a CDN for content that's the same for all users.
6. Talk about edge cases like cache stampedes and hot keys, especially for data that can suddenly become very popular.
Following these steps, each component you add is tied to a specific problem you identified earlier in the discussion. This makes the answer stronger than a list of technologies.
20 quizzes