Scaling Out: Load Balancers, Caches and CDNs
Most system design answers rest on the same four ideas: spread load across many machines, keep those machines stateless, put frequently read data close to the reader, and avoid doing work twice. This chapter covers the components that implement them and the questions interviewers ask about each.
1. Vertical and horizontal scaling
Vertical scaling means a bigger machine: more CPU, memory and faster disks. It is simple, needs no code change, and is bounded by the largest machine you can buy and by the fact that one machine is one failure.
Horizontal scaling means more machines behind a front door. It has no hard ceiling and gives redundancy, but it forces you to deal with distribution: where state lives, how requests are routed, and how machines agree on anything.
A fair answer to "how would you scale this?" starts vertical (cheap, quick), then moves horizontal when you hit the limit or need availability. Most real systems do both.
2. Stateless application servers
Horizontal scaling of the application tier works only if any server can handle any request. That means no user state in server memory. Move it out:
- session data to a shared store such as a cache or database, or into a signed token carried by the client;
- uploaded files to object storage;
- background work to a queue.
With stateless servers, you can add or remove machines at will, replace a failed one without losing sessions, and deploy gradually. The state that remains lives in a small number of purpose-built stores, which are then the hard part of the design.
<!--fig:lb-->3. Load balancers
A load balancer accepts client connections and distributes them over a pool of servers. It also performs health checks, removing servers that fail and returning them when they recover.
Layer 4 versus layer 7. A layer 4 balancer routes on IP and port and does not read the request, so it is fast and protocol-agnostic. A layer 7 balancer understands HTTP: it can route by path or header, terminate TLS, rewrite requests and apply rate limits, at the cost of more processing.
Common algorithms:
| Algorithm | How it picks | Best when |
|---|---|---|
| Round robin | Next server in turn | Requests are similar in cost |
| Weighted round robin | Turns proportional to capacity | Servers differ in size |
| Least connections | Server with the fewest active connections | Request time varies widely |
| Hash of client or key | Same input maps to the same server | You need affinity, such as local caches |
Sticky sessions pin a client to one server. They are a workaround for stateful servers and should be a last resort, since they unbalance load and lose the session when the server dies.
The balancer is itself a single point of failure. Real deployments run several, behind DNS or a virtual IP with failover, or use a managed service that does so for you. Mention this unprompted.
4. Caching
A cache stores the result of expensive work so the next request is cheap. It works because access is skewed: a small fraction of the data serves most of the reads.
Where caches live
- Browser and client caches, controlled by HTTP headers.
- CDN edge caches, near the user (next section).
- Application-level caches such as an in-memory store shared by all servers.
- Database caches, such as the buffer pool.
Patterns for reading and writing
Cache-aside (lazy loading). The application checks the cache, on a miss reads the database, then fills the cache. It is the most common pattern, tolerates cache failure, and caches only what is requested. Its weakness is that a cache miss on a cold or recently flushed cache sends every request to the database at once.
<!--fig:cacheaside-->Read-through. The cache itself loads data on a miss. The application sees one interface.
Write-through. Writes go to the cache and the database together, so reads are fresh, but each write is slower.
Write-back (write-behind). Writes go to the cache and are flushed to the database later. Writes are fast, but data can be lost if the cache dies before flushing, so use it only where that risk is acceptable.
Invalidation
Phil Karlton's remark that cache invalidation is one of the two hard problems in computer science survives because it is true. Options:
- Time to live (TTL): the entry expires after a set time. Simple, and it bounds staleness.
- Explicit invalidation: delete or update the entry when the source changes. Fresher, but you must find every place that writes.
- Versioned keys: include a version in the key so a change creates a new entry and the old ones age out.
Say which staleness the product tolerates. A profile photo can be minutes stale. An account balance cannot.
Eviction
When the cache is full, an eviction policy chooses what to drop. Least recently used (LRU) is the default answer, least frequently used (LFU) protects genuinely popular items, and first in, first out is rarely best. Be ready to implement LRU with a hash map and a doubly linked list, which gives constant-time get and put.
Cache failures to know by name
- Cache stampede (thundering herd). A popular key expires and thousands of requests miss together. Mitigate with request coalescing (one request refills while others wait), randomised TTLs, or refreshing before expiry.
- Cache penetration. Requests for keys that do not exist always miss and hit the database. Cache the negative result briefly, or filter with a Bloom filter.
- Hot key. One key receives so much traffic that a single cache node overloads. Replicate the key across nodes or add a small local cache in each application server.
- Cold start. A freshly started cache is empty. Warm it before sending it traffic.
How big should the cache be?
A rule of thumb used in interviews is the 80/20 rule: if 20 percent of the data serves 80 percent of the reads, caching that 20 percent absorbs most of the load. Size the cache from your estimate: if a day's hot set is 20 GB, a 32 GB cache is plausible, and you measure the hit rate to confirm.
5. Content delivery networks
A CDN is a network of servers placed in many regions that cache content close to users. The first request for an object goes to the origin, and subsequent requests nearby are served from the edge.
<!--fig:cdn-->It cuts latency, reduces load on the origin, and absorbs traffic spikes and some attacks. It suits static assets (images, scripts, video segments) and increasingly cacheable API responses.
Details worth mentioning:
- Cache keys and headers.
Cache-ControlandETagtell the edge how long to keep an object and how to revalidate it. - Invalidation. Purge by URL or tag, or use versioned file names so the old object is simply never requested again. Versioned names are the more reliable option.
- Push versus pull. In pull mode the edge fetches from the origin on demand. In push mode you upload content in advance, useful for large, predictable files.
- Private content. Use signed URLs or cookies so that only authorised users can fetch from the edge.
Potential deep dives
Deep dive 1: How do you keep application servers stateless?
The challenge. A user logs in on one request and the next request lands on a different server. The second server must still know who they are.
Weak: sticky sessions. The load balancer pins each user to one server, which keeps the session in its own memory. It works until that server dies or is replaced during a deployment, and then every pinned user is logged out. It also unbalances load, because heavy users stay on one machine, and it makes scaling in and out disruptive.
Solid: a shared session store. Store the session in a fast shared store such as a replicated in-memory cache, keyed by a session identifier in a cookie. Any server can look it up. Servers become interchangeable, and a server loss costs nothing.
Excellent: stateless tokens where possible, and a store where not. Carry the identity in a signed, short-lived token that the client sends with each request, so no lookup is needed at all, and keep a small server-side store only for things that must be revocable or large. State the trade-off: a signed token cannot be revoked before it expires unless you add a denylist, so keep lifetimes short and use refresh tokens. Keep everything else, such as uploads and in-progress work, out of server memory too, in object storage and queues.
Deep dive 2: How should the load balancer distribute traffic?
The challenge. Requests differ in cost, servers differ in capacity, and servers fail.
Weak: round robin and nothing else. It ignores that one request may take a thousand times longer than another, so a server stuck with several slow requests gets overloaded while others idle. It also keeps sending traffic to a server that has failed until someone notices.
Solid: least connections with health checks. Send each request to the server with the fewest active connections, which adapts to uneven request cost. Probe each server periodically and remove those that fail, returning them after they recover. Use weights when servers differ in size.
Excellent: layered, observable and failure-aware. Use a layer 4 balancer in front for raw speed and a layer 7 tier for routing by path, header or tenant, with TLS terminated at the edge. Make health checks meaningful: a readiness check that verifies the server can reach its dependencies, not only that the process is alive. Add connection draining so that a server being removed finishes its in-flight requests, and slow start so that a new server ramps up gradually. Use outlier detection to eject a server whose error rate rises. Run the balancers themselves redundantly, behind DNS or a virtual address with failover, and watch the distribution of load per server.
Deep dive 3: What do you cache, and how do you keep it correct?
The challenge. A cache speeds reads and creates a second copy of the data that can be wrong.
Weak: cache everything with a long TTL. It hides bugs, serves stale data for hours, and wastes memory on items that are read once.
Solid: cache hot, expensive, shared data with a TTL, using cache-aside. Pick items by access frequency and cost to compute, set a TTL based on how stale the product tolerates, and delete the entry when the underlying data changes.
Excellent: reason about hit rate, invalidation races and failure. Estimate the hit rate from the skew of access and size the cache to the hot set. State the invalidation strategy per data type: TTL only for tolerant data, delete-on-write for data that must be fresher, and no caching for data that must be exact. Explain the stale-fill race (a reader loads old data, a writer updates and invalidates, then the reader fills the cache with the old value), and bound it with a short TTL as a safety net. Protect against stampedes with request coalescing and jittered TTLs, hot keys with replication or a local cache, and cache penetration with negative caching. Make sure the database can survive a cold or failed cache, because otherwise the cache is a hidden single point of failure.
Deep dive 4: How do you serve users around the world?
Weak: one region. Users on the other side of the planet pay hundreds of milliseconds on every request, and a regional outage takes everyone down.
Solid: a CDN for static content and a second region as a standby. Static assets and cacheable responses come from the edge. A standby region is ready for failover.
Excellent: place data near users, and be honest about writes. Serve static and cacheable content from the edge, replicate read-heavy data to regions near users, and route each user to the nearest healthy region. Decide where writes go: a single leader region with replicas elsewhere (simple, with cross-region write latency), partitioning users by home region so most writes are local, or multi-region writes with conflict resolution (powerful and complex). Say that cross-region synchronous replication makes every write wait a round trip across continents, so most designs replicate asynchronously and accept a small window of data loss on regional failure.
What is expected at each level
Mid-level. You separate stateless servers from state, add a load balancer and a cache, and explain cache-aside. You know what a CDN does.
Senior. You compare load-balancing algorithms and health checks, handle sessions properly, reason about invalidation and stampedes, and size caches from numbers. You make the database survive cache failure.
Staff. You reason about the layers together: where each request is served, how failures cascade across layers, how multi-region design affects consistency and cost, and how you would operate and observe it all.
Interview questions and model answers
Q: A page is slow. Where do you add a cache? I measure first to find where time goes. If the cost is repeated database reads, I cache the result in a shared in-memory store using cache-aside with a TTL. If the cost is static assets crossing the world, I put them on a CDN. I pick the layer closest to the user that the data's freshness allows.
Q: What happens when your cache goes down? With cache-aside the application still works, reading from the database, but the database must be able to take the full read load or it will fall over. I would protect it with rate limiting, a smaller local cache and a replicated cache cluster so a single node loss does not empty everything.
Q: Why not use sticky sessions? They tie a client to one server, which unbalances load and loses the session on failure. Making servers stateless and storing sessions in a shared store gives the same user experience without those problems.
Q: How do you keep a cache consistent with the database? I cannot get perfect consistency cheaply. I choose between a TTL that bounds staleness and explicit invalidation on writes, and I decide per data type how stale a reader may be. For data that must never be stale, I do not cache it, or I read from the primary.
Common mistakes
- Adding a cache without saying what invalidates it.
- Treating the load balancer as infinitely reliable.
- Keeping session state in application memory.
- Forgetting that the database must survive a cold cache.
- Caching personalised responses at a shared edge.