AI System Design
← All chapters
Chapter 1 7 min read

Scale from Zero to Millions of Users

Growing a web product from one box to millions of users is not one big leap but a sequence of small, well-understood moves: split the database out, add redundancy, cache aggressively, push static bytes to the edge, make web servers stateless, and decouple work with queues. Each step removes the current bottleneck or single point of failure and exposes the next one. This chapter is the vocabulary every later design builds on.

Architecture at a glance
  1. Client (browser / mobile app)
  2. GeoDNS
  3. CDN for static assets
  4. Load balancer
  5. Stateless web tier
  6. Cache (Redis / Memcached)
  7. Primary + replica databases (sharded)
  8. Message queue + async workers

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.
Figure 1Cache-aside read, then a write
Cache-aside read, then a writeClientWeb serverCacheReplica DBPrimary DB1. GET /profile/422. get user:423. miss4. SELECT user 425. row (may lag the primary)6. set user:42, TTL 10 min7. 200 OK8. PUT /profile/42 (new bio)9. UPDATE user 4210. replicate change (async)11. delete user:42 (invalidate)
Reads try the cache, fall back to a replica, and repopulate the cache; writes go only to the primary, which streams changes to replicas asynchronously. Deleting the cache entry after the write keeps the next read from serving the old value for a whole TTL.

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.

Figure 2Two data centers behind GeoDNS
Two data centers behind GeoDNSDATA CENTER 1DATA CENTER 2resolveresolveasync replicationUS usersGeoDNSanswers with nearest siteEU usersLoad balancerWeb serversstateless, autoscaledDatabasescache + primary/replicasIf one site failsDNS shifts users to the otherLoad balancerWeb serversstateless, autoscaledDatabasescache + primary/replicas
GeoDNS sends each user to the nearest site, and each site runs a complete stateless stack. Data is replicated asynchronously between sites, so when one site fails DNS can shift everyone to the survivor without losing more than the replication lag.

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.

Figure 3The evolved single-region architecture
The evolved single-region architectureASYNC PROCESSINGlookupstaticHTTPSreadfirstsession datawritesreadsreplicateenqueue jobconsumeDNSdomain → LB / CDNCDNimages, JS, CSSMessage queuedurable job bufferWorkersthumbnails, email, …Usersweb + mobileLoad balancerhealth checksWeb serversstateless, autoscaledCache clusterRedis / MemcachedSession storeshared, not on web hostsPrimary DBall writesRead replicasmost reads
Every box removes one bottleneck: the CDN takes static bytes, the cache absorbs hot reads, replicas take the remaining reads, and the queue moves slow work off the request path. Because sessions live in a shared store, any web server can serve any request.

Sharding the data tier

Eventually the write volume or dataset outgrows one primary. Sharding splits data across several databases with the same schema, chosen by a sharding key such as user_id % number_of_shards. A good key spreads load evenly and keeps the data a query needs on a single shard. Sharding brings real costs, which is why it is usually the last step rather than the first.

  • Resharding: shards fill unevenly or the count must grow; consistent hashing limits how much data must move.
  • Celebrity (hotspot) problem: one key, like a famous account, can overload a shard; give such keys dedicated shards or split them further.
  • Joins and denormalization: cross-shard joins are expensive, so data is often duplicated so each query hits one shard.
  • Moving some functionality to NoSQL stores can relieve relational shards entirely.

Key numbers

Seconds per day
86,400 (round to 100,000)
10M DAU x 20 req/day
~2,300 QPS avg, ~11,500 peak at 5x
Hot 20% of 50M x 1 KB objects
~10 GB of cache
Typical read:write ratio
10:1 or higher
Cache vs DB read latency
sub-millisecond vs several ms

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