AlgoMaster Logo

Scaling Read Traffic

High Priority19 min readUpdated September 29, 2026
AI Mock Interview

Practice this topic in a realistic system design interview

Listen to this chapter
Unlock Audio

Scaling Read Traffic

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.

Where This Pattern Shows Up

Let's start with the kinds of problems where reads are much more common than writes.

ProblemWrittenRead
URL shortener (Bitly)A short link is created onceIt can be clicked millions of times
Design Twitter or InstagramA post is written onceIt's read by every follower who opens their feed
Design YouTubeA video's title, description, and view count change rarelyThey're loaded every time someone opens the video page
E-commerce (Amazon)Product details are updated occasionallyProduct pages are viewed far more often
Ticket booking (Ticketmaster)A ticket is bought onceThousands 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.

1. Why Reads Become a Problem

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.

2. Start With the Database Itself

Before adding new components, it's worth checking whether the database is being used efficiently.

2.1 Indexes

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.

2.2 The Queries Themselves

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.

ProblemFix
Full table scanAdd an index on the lookup column
N+1 queriesUse a JOIN or batch fetch with IN (...)
Large result setAdd LIMIT and paginate
Expensive aggregation on every readPrecompute 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.

3. Read Replicas

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.

4. Replication Lag

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.

Change not yetcopied to the replica"Did my update fail?"Update profile picture200 OKReload the pageOld profile pictureUserPrimaryReplicaUserPrimaryReplica
6 / 6
algomaster.io

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.

FixHow It WorksTrade-off
Primary after a writeRoute a user's reads to the primary for a few seconds after they writeAdds some load to the primary
Own data from primaryA user's own profile, settings, and orders always read from the primaryMore primary load for data-heavy users
Lag monitoringRemove replicas from rotation when they fall too far behindFewer 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.

5. Caching

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.

5.1 Cache-Aside

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.

alt[Cache Hit][Cache Miss]GET /a1B9xQGET link:a1B9xQlong URL302 RedirectnullSELECT long_url WHERE short_code = 'a1B9xQ'long URLSET link:a1B9xQ (with TTL)302 RedirectClientApp ServerRedis CacheDatabase
9 / 9
algomaster.io

5.2 The Cache Hit Ratio

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.

5.3 What to Cache

Good candidates are read often, change rarely, and can tolerate a little staleness.

Good Cache CandidatesPoor Cache Candidates
Short link to long URL mappingsBank account balances
Product detailsPrices that must be quoted exactly at checkout
Posts and profilesOne-time authentication tokens
Configuration dataData that is rarely read

6. Keeping the Cache Up to Date

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.

6.1 Time to Live (TTL)

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.

6.2 Delete on Write

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.

6.3 Use Both

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:

  • Write-through updates the cache in the same write path as the database. The cache stays fresher, but every write gets slower and you have to decide what happens when the database write succeeds and the cache write fails.
  • Event-based invalidation publishes a change event to a queue like Kafka, and each service that caches the data deletes its own copy. This decouples the writer from the caches, which helps when many services cache the same data.
StrategyStalenessWrite CostComplexityGood Fit
TTL onlyUp to the TTLNoneLowData where a few minutes of staleness is fine
Delete on write + TTLVery short, bounded by TTL if a delete is missedOne extra deleteLowMost applications
Write-throughVery shortHigher, cache in the write pathMediumSmall read models that must stay fresh
Event-basedShort, depends on the queueLow for the writerHighMany services caching the same data

7. Cache Stampedes

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.

TechniqueHow It WorksTrade-off
Lock on rebuildOne request reloads the entry; the rest wait or get the old valueWaiting requests add a little latency
Background refreshPopular entries are reloaded before they expireExtra work for entries that are no longer read
TTL jitterRandom variation spreads expirations over timeDoes not help a single hot key on its own

8. Hot Keys

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.

9. CDNs and Edge Caching

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.

ContentExample Cache-ControlMeaning
Versioned images, CSS, JSpublic, max-age=31536000, immutableCache for a year; a new version gets a new URL
Public product page or articlepublic, max-age=60, s-maxage=300Browser reuses it for 60 s, the CDN for 300 s
User-specific dataprivate, no-storeDo 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.

10. Precomputing Results

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.

11. Putting the Layers Together

Let's now look at how all of these layers work together.

  1. A read request first reaches the CDN. If the CDN has the response, it returns it right away.
  2. If not, the request reaches your application server, which may check its own in-memory cache.
  3. If the data isn't there, it checks the distributed cache.
  4. And only if that misses does the request reach a read replica, or the primary database.

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.

LayerHow Stale It Can BeHow It Gets Updated
BrowserUp to max-ageExpires, or a new versioned URL
CDNUp to s-maxageExpires, or an explicit purge
App server memoryA few secondsShort TTL
Distributed cacheUp to the TTL, usually much lessDelete on write, plus TTL
Read replicaReplication lag, usually millisecondsReplication stream from the primary

12. Discussing High Read Traffic in an Interview

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.

QuestionWhy 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.

13. Key Takeaways

  1. Start with efficient queries and indexes. They reduce database load without adding any new infrastructure.
  2. Cache frequently read data. The cache hit ratio decides how much load reaches the database: at 95%, 35,000 reads per second becomes 1,750.
  3. Keep the cache up to date by deleting entries on writes and setting a TTL in case a delete is missed.
  4. Add replicas for the reads that miss the cache, and handle replication lag for users who just wrote something.
  5. Use a CDN for content that's the same for every user, including public pages and API responses, not just images.
  6. Precompute expensive reads at write time, like feeds and counters, when reads are much more common than writes.
  7. Plan for popular data. Stampedes and hot keys can overload a cache server or the database even when the rest of the caching design works.
  8. At each step, think about how stale the data can be, and how each layer gets updated.

Quiz

Scaling Read Traffic Quiz

20 quizzes