Problem & requirements
The starting point is a product that runs everything on a single machine: web server, application code, database and static files. That is the right call on day one because it is cheap and simple to reason about, but it has no redundancy and a hard ceiling on capacity. The goal of the chapter is to evolve that box into a system that stays available when individual machines die, scales roughly linearly as traffic grows, and keeps latency low for users who may be on the other side of the planet.
Implicit requirements shape every decision: reads typically dominate writes, traffic is spiky (product launches, time-of-day peaks), and the team wants to add capacity without rewriting the application. Keep those three facts in mind and most of the moves below follow naturally.
- Availability: survive the loss of any single server, and eventually of a whole data center.
- Scalability: add capacity by adding machines rather than replacing them.
- Performance: serve hot reads from memory and static content from the edge.
- Operability: observe, deploy and recover without heroics.
Back-of-the-envelope
Suppose the product reaches 10 million daily active users, each making about 20 requests a day. That is 200 million requests per day; dividing by 86,400 seconds gives roughly 2,300 requests per second on average. Traffic is uneven, so plan for a peak of about 5x, or roughly 11,500 QPS. If one application server comfortably handles around 1,000 dynamic requests per second, the web tier needs a dozen machines at peak plus headroom for a failure, which is exactly why a load-balanced, horizontally scaled tier becomes necessary.
For the cache, imagine 50 million objects of about 1 KB each. Following the usual 80/20 skew, caching the hottest 20% means 10 million objects times 1 KB, about 10 GB, which fits in the RAM of a single cache node. If 10% of requests are writes, the database sees about 230 writes per second on average and 1,150 at peak: fine for one primary, while reads are offloaded to replicas and the cache.
From one server to separate tiers
The first split is to move the database onto its own machine so the web and data tiers can scale independently. This raises the classic choice between relational and non-relational stores. Relational databases give you joins, transactions and decades of tooling, and they are the sensible default. NoSQL stores (key-value, wide-column, document, graph) earn their place when you need very low latency on simple lookups, have unstructured or rapidly evolving data, or must store volumes that are awkward to shard in SQL.
Next comes the scaling choice. Vertical scaling (a bigger box) is operationally trivial but hits a hardware ceiling and leaves a single point of failure. Horizontal scaling (more boxes) has no practical ceiling and gives redundancy, at the cost of needing a load balancer in front. The load balancer exposes one public IP, talks to web servers over private addresses, health-checks them, and routes around dead instances, so adding capacity becomes a matter of registering another server.
Database replication
With the web tier redundant, the database becomes the weakest link. Primary-replica replication keeps one primary that accepts all writes and several replicas that receive a copy of its change log and serve reads. Because most workloads are read-heavy, this multiplies read capacity, improves availability, and lets replicas live in different racks or regions so a local disaster does not destroy all copies of the data.
Failure handling is where the subtlety lives. If a replica dies, reads shift temporarily to others or to the primary. If the primary dies, a replica is promoted, but asynchronous replication means it may be missing the last few writes, which must be reconciled or accepted as lost. Replication lag also means a user may write and then immediately read stale data from a replica; a common fix is to route a user's reads to the primary for a short window after they write (read-your-writes).
Caching and the CDN
A cache is a fast in-memory layer that absorbs repeated reads. The most common pattern is cache-aside: the application checks the cache, falls back to the database on a miss, then populates the cache. Key decisions are the expiration policy (too short and you hammer the database, too long and data goes stale), the eviction policy when memory fills (LRU is the usual default, LFU or FIFO in special cases), and consistency, because updating the database and the cache is not atomic. A single cache node is also a single point of failure, so production caches run as a replicated cluster with some memory over-provisioned.
A CDN applies the same idea geographically to static assets: images, video, CSS and JavaScript are served from edge locations near users. You pay per byte transferred, so rarely requested files are better served from origin. Set sensible TTLs, version asset URLs so you can invalidate by publishing a new name, and design a fallback so clients can reach the origin if the CDN has an outage.
- Read-through and write-through caches hide the cache behind a library or the store itself.
- Write-behind batches writes for speed but risks loss if the cache crashes.
- Cache only data that is read far more often than it changes.
Stateless web tier and multiple data centers
Horizontal scaling only works cleanly if any server can handle any request. When session data lives in a server's memory, the load balancer must use sticky sessions, which complicates adding, removing and failing over servers. The fix is to move state into a shared store (Redis, a NoSQL table or the database) so web servers become stateless and can be autoscaled freely based on load.
To survive a regional outage and to cut latency, the system can run in multiple data centers. GeoDNS resolves users to the nearest healthy site, and if one site fails all traffic is shifted to the survivors. This is harder than it sounds: data must be replicated between sites, caches warmed, and deployments kept identical, and you must test failover regularly or it will not work when needed.
Message queues, observability and automation
A message queue is a durable buffer that decouples producers from consumers. A web server can accept an expensive job, such as resizing a photo, publish it, and return immediately while a pool of workers drains the queue. Producers and consumers scale and fail independently: if workers fall behind, the queue grows and you add workers; if they crash, messages wait instead of being lost. This decoupling is one of the most reusable ideas in system design.
Once there are dozens of machines, you need logging aggregated in one place, metrics at host, aggregate and business levels (CPU, p99 latency, daily active users, revenue), and automation in the form of continuous integration and deployment. These are not luxuries; without them you cannot find the bottleneck you are supposed to scale next.
Key numbers
Key terms
- Vertical scaling
- Adding CPU, RAM or disk to an existing machine; simple but capped and non-redundant.
- Horizontal scaling
- Adding more machines to a pool so capacity and redundancy grow together.
- Load balancer
- A component that spreads incoming traffic across healthy backend servers and hides them behind one address.
- Primary-replica replication
- One node takes writes and streams changes to read-only copies that serve reads.
- Cache-aside
- A pattern where the application reads from cache first and loads from the database on a miss.
- CDN
- A geographically distributed network of servers that caches static content close to users.
- Stateless tier
- Servers that keep no per-user state locally, so any instance can serve any request.
- GeoDNS
- DNS that answers with different IPs depending on the requester's location.
- Sharding
- Horizontal partitioning of data across databases by a key so each holds a subset.
Common mistakes
- Jumping straight to sharding and microservices before replication, caching and a stateless tier have been exhausted.
- Adding a single cache node and forgetting it is now a single point of failure with a cold-start problem.
- Ignoring replication lag, so users do not see their own just-written data.
- Choosing a sharding key with a skewed distribution, creating permanent hotspots.
- Keeping sessions in web-server memory and then wondering why autoscaling breaks logins.
- Proposing multi-region active-active without addressing how data is replicated and conflicts resolved.
Further study
- Designing Data-Intensive Applications (Martin Kleppmann)
- Scaling Memcache at Facebook (NSDI 2013)
- The Twelve-Factor App methodology
- Google SRE book