Course Content
System Design Interview
31 sections · 71 lessons
Scaling the web tier: from one server to a stateless pool
Every large system started as one machine running everything, and understanding that machine completely is the fastest route to understanding why the large one looks the way it does.
This section grows one system from a single box to millions of users, one step at a time. Each step is forced by a specific failure of the step before it, and each one costs something. This lesson covers the first four steps: the single server, splitting out the database, choosing between a bigger machine and more machines, and putting a load balancer in front of a stateless web tier. The next lesson scales the data itself.
The running example
Throughout this section we grow one invented product: Trailmix, a route-sharing app for cyclists. Users record a ride, upload it with a few photos, and browse routes posted by people they follow. The company, the users, and every number below are made up, but the arithmetic is real and you should follow it.
Trailmix launches with about 1,000 daily active users — people who open the app on a given day. Each opens it five times and each session makes ten requests. That is 50 requests per user per day, so 50,000 requests a day, which is 50,000 ÷ 86,400 ≈ 0.6 requests per second on average. Even at a three-times evening peak that is under 2 requests per second.
One machine is not a compromise here. It is comfortably over-provisioned.
Where the request actually goes
1. Name resolution. The client has a domain name and needs an internet protocol (IP) address. It asks the Domain Name System (DNS) — a global lookup service mapping names to addresses. A cold lookup crosses several servers and takes roughly 20–100 ms. It is then cached by the browser, the operating system, and the network resolver for the record's time-to-live (TTL), so almost every later request skips it entirely.
2. Connection setup. The client opens a Transmission Control Protocol (TCP) connection — one round trip — then negotiates Transport Layer Security (TLS) for encryption, historically two more round trips and one in modern versions. On a 40 ms round trip that is 80–120 ms before a single byte of your request has been sent. This is why connection reuse matters more than most application-level optimisation.
3. The application. The web server accepts the request, the application code runs, and it queries the database. An indexed lookup against data already in memory takes well under a millisecond.
4. The response. JavaScript Object Notation (JSON) goes back over the same connection.
What one box can actually do
A modern eight-core machine with 32 GB of memory, running a web server, an application process, and a relational database, will serve on the order of a few thousand trivial requests per second and a few hundred that involve real database work. Trailmix needs two.
This matters because candidates reflexively reach for a distributed architecture on a system that would run on a laptop. Plenty of real, profitable businesses run on one large machine and a backup. Saying so out loud in an interview is a credibility signal, not a weakness — as long as you follow it with the thing that actually breaks.
What breaks first
Not the central processing unit (CPU). What breaks first is that there is one of everything. The machine reboots for a kernel update and Trailmix is down. A bad deploy takes the database with it. A runaway query starves the web server of memory. There is no way to add a second web server, because the database lives inside the first one.
One machine handles far more than most candidates assume — Trailmix's launch traffic uses about 1% of it. The first thing that forces a change is not capacity. It is that a single box is a single point of failure and cannot be grown in parts.
Separating the database
The first split is always the same: move the database onto its own machine. It happens first because it unblocks everything else.
There are three reasons for this split and not another, in order of how much they matter.
1. It makes the web tier growable. As long as the database lives inside the web server, adding a second web server means adding a second, separate copy of your data. Splitting the database is the precondition for every later step in this section.
2. The two workloads want different machines. A database wants memory (to keep hot data cached) and fast disk. An application process wants CPU. Sharing 32 GB between a database buffer pool and a fleet of application workers means both are short, and the failure is silent: the database quietly starts reading from disk instead of memory, and a query that took 0.5 ms starts taking 5 ms.
3. They fail independently. A memory leak in your application no longer takes down the store of record.
Trailmix now has two machines: a web/app server and a database server, talking over the internal network. The cost is one network round trip per query — roughly 0.5 ms inside a datacentre. A page that runs eight queries now pays 4 ms it did not pay before. That is a real cost and it is worth it.
Relational or non-relational
Once the database has its own machine, the next question is what kind of database it should be.
Relational database: data in tables of rows with a fixed schema, related by keys, queried with Structured Query Language (SQL). It supports joins (combining rows from several tables in one query) and transactions with ACID guarantees — Atomicity, Consistency, Isolation, and Durability, meaning a group of writes either all happen or none do, and survives a crash.
Non-relational ("NoSQL"): a family, not one thing — key-value stores, document stores, wide-column stores, and graph stores. The Section 5 lesson Databases: choosing and justifying covers all of them properly. What they share is dropping some combination of joins, schemas, and multi-row transactions in exchange for easier horizontal growth and, sometimes, a much better fit to one access pattern.
The honest answer
For the core data of most interview systems, choose relational. Trailmix's data is users, routes, and follows — entities with relationships, queried in combinations you have not thought of yet. "Show me routes posted in the last week by people I follow, ordered by distance" is one SQL statement against three tables. In a document store, it is either a denormalised copy of the follow graph inside every route or three round trips and a merge in application code.
Reach for a non-relational store when you have a specific reason: an access pattern that is purely key-based (Section 10, URL shortener), a write volume no single leader can absorb (Section 14, chat), or data with no useful relations, like metrics (Section 22).
The sentence to say in an interview: "Relational for the core entities because the access pattern involves relationships and I want transactions; I would move X to a key-value store because its access pattern is a single key lookup at high volume."
Vertical versus horizontal scaling
With the database on its own machine, the web tier can grow. There are two ways to handle more traffic: a bigger machine, or more machines. They are not interchangeable, and the point at which you switch is a real decision with a real cost.
Vertical scaling (scale up): replace the machine with a larger one — more cores, more memory, faster disk. No code changes.
Horizontal scaling (scale out): add more machines and spread the work across them. Code changes, usually significant ones.
Why vertical scaling is the right first move
It is the cheapest thing that works. Trailmix grows to 100,000 daily active users: 5 million requests a day, about 58 per second average and perhaps 175 at peak. Doubling the web server from 8 to 16 cores handles that with no engineering effort at all. A week of engineering time costs more than a year of the larger instance.
Candidates who leap to a sharded, queue-driven architecture at 175 requests per second are demonstrating that they cannot size a problem.
Why it stops working
There is a ceiling. The largest instances cloud providers rent have a few hundred cores and single-digit terabytes of memory. That is a hard stop, and it is closer than it sounds once one process must hold everything.
The price curve bends. Going from 8 to 16 cores roughly doubles the price. Going from 64 to 128 cores typically costs more than double, because you are now buying the top of the range where there is less competition. (Indicative — check current pricing, it moves.)
It does not remove the single point of failure. One enormous machine is still one machine. It reboots, and you are down.
Upgrades mean downtime. Resizing an instance means stopping it.
What horizontal scaling costs you
This is the part candidates skip, and it is the interesting part.
Statelessness. Anything held in one server's memory — a session, an in-process cache, an uploaded file waiting to be processed — becomes wrong the moment there are two servers. The load balancer below only works once this is dealt with.
Coordination. Two servers that both want to send one email, or both want to assign the next order number, now need a shared source of truth. Section 9 is an entire section on generating unique identifiers once one machine is no longer doing it.
Partial failure. On one machine, things are up or down. With twenty machines, three can be up but slow, one can be up but serving stale data, and one can be unreachable from half the fleet but not the other half. Every reliability pattern in Section 5 (Reliability patterns) exists because of this.
Debugging gets harder. A request now touches five machines. You need distributed tracing (Observability and deployment) to follow it at all.
The practical rule
Scale vertically until one of three things is true: you are near the top of the instance range, the cost curve has bent badly, or you need more than one machine for availability rather than capacity. In practice the third arrives first, which is why the next step is a load balancer rather than a bigger box.
Load balancers
A load balancer is a machine that receives every client request and forwards it to one of several identical web servers. It converts a single point of failure into a pool.
Clients resolve trailmix.app to the load balancer's public IP address. The web servers sit behind it on private addresses, unreachable from the internet. The balancer keeps a list of healthy servers, picks one per request, forwards it, and returns the response.
Health checks are what make it useful. The balancer requests a known path — say /healthz — from each server every few seconds. Miss the threshold and the server is pulled from rotation; pass again and it returns. Now a crashed server means a few seconds of errors instead of an outage, and a deploy means draining one server at a time while the others serve.
Load balancing and proxies, in Section 5, covers the routing algorithms (round robin, least connections, hash-based) and the difference between layer 4 and layer 7 balancing. For now, round robin is fine.
Sizing the pool
Each Trailmix web server handles about 200 requests per second of real work. At a peak of 175 per second, one is enough — until it dies. Provision N+1: two servers, either of which can carry the whole peak. At 1 million daily active users the peak is around 1,700 requests per second, so nine servers of capacity plus one spare, and the pool grows linearly from there.
The state problem
This is the real content of the load-balancer step. Here is a naive login flow that works perfectly on one server and breaks on two.
A user signs in. The server generates a session identifier, stores {session_id → user_id, expires_at} in its own memory, and returns the identifier in a cookie. Every later request presents the cookie, and the server looks it up in memory.
With two servers behind a round-robin balancer, request one lands on server A, which creates the session. Request two lands on server B, which has never heard of that session, and the user is logged out. In practice this shows up as users being randomly signed out on about half their requests.
Three fixes, and which to pick
Sticky sessions. The balancer pins each client to one server, usually via a cookie. It works, and it is the wrong answer for three reasons: when that server dies its users all lose their sessions; load becomes uneven, because the pinning does not know how heavy each user is; and you cannot drain a server for deployment without disrupting its users.
Shared session store. Sessions move to a store every web server can reach — a key-value store such as Redis, or a database table. One extra lookup per request (roughly 0.5 ms in the same datacentre). Any server can serve any user. Now the web tier is genuinely stateless.
Signed tokens. Put the session data in a token the client holds, cryptographically signed so the server can verify it without storing anything. No lookup at all. The trade-off is revocation: a token is valid until it expires, so signing a user out immediately requires a denylist, which reintroduces the shared store you were avoiding.
Recommendation: shared session store for most systems, because a 0.5 ms lookup buys immediate revocation and one less thing to reason about. Signed tokens when the lookup cost genuinely matters and short expiry is acceptable.