From Zero to Millions: Scaling a Web Service

The classic journey from one server to millions of users: vertical and horizontal scaling, load balancers, read replicas, CDNs, and sharding — with the napkin math to back it up.

Beginner · 15 min read

Why this matters

Almost every system design interview starts the same way: "Design a web service that goes from zero users to a million." The interviewer isn't really asking about a million — they're asking whether you know the stages. Every large system you admire — a video platform, a chat app, a food delivery service — walked the same path, solving one bottleneck at a time. This lesson walks it with you, step by step. Learn this journey once and you'll recognize it everywhere.

Video: System Design Introduction For Interview. — Tushar Roy - Coding Made Simple
Starts from a single server and grows into the full scaling toolkit, interview-style — the same zero-to-millions journey this lesson walks.

Picture a food truck

This lesson's running analogy: you run a food truck. One truck, one cook, one window for taking orders. That's your web service on day one.

Mapping, stated plainly:

Day one is glorious. Customers trickle in. The cook handles it. Everything lives in the one truck. The problem is what happens when the food blog writes you up and five hundred people show up at once. That moment has a name: scaling. And there are two fundamentally different ways to do it.

Video: Horizontal vs Vertical Scaling: How to Handle Million Users | System Design Ep 2 — Vishal Kumar
Uses a restaurant analogy to make horizontal and vertical scaling concrete, building the mental model this stage needs.

Stage 0: One truck (and its limits)

Your first server does it all: it serves the web pages, runs your application code, and holds the database on the same disk. For a hobby project this is honestly the right call — one machine, one bill, one thing to debug at 2 a.m.

The first move when traffic grows is the obvious one: scale vertically. Get a bigger truck. More RAM, a faster CPU, a bigger disk. In the analogy, that's a bigger grill, a wider prep counter, a second register.

Vertical scaling has real charms: nothing in your code changes, your recipe book stays in one place, and there's no distributed anything to keep you up at night. But it has a hard ceiling. There's only so big a grill can get, and the price curve is cruel — the jump from a good machine to a great one costs far more than the performance you gain. Worse, you still have exactly one truck. If the engine dies, lunch is cancelled. Vertical scaling has no answer for failure.

So the industry converged on the other direction.

Video: Vertical vs Horizontal Scaling: Bigger Box or More Boxes? — The Hot Path
Shows exactly where one server hits its limits and why vertical scaling eventually stops working — Stage 0's core lesson.

Stage 1: Clone the truck (horizontal scaling)

Horizontal scaling means: instead of one bigger truck, run more trucks. Ten ordinary trucks side by side can serve more customers than one enormous one, and if one breaks down, the other nine keep cooking. You trade one super-machine for a fleet of cheap ones.

But fleets create a new problem. When a customer walks up, which truck takes their order? You need someone at the front of the line pointing people to the shortest queue. That's the load balancer — the friendly order-taker with the headset who waves each customer to a free window.

Here's the whole idea in one picture:

flowchart LR
    U[Customers]:::client --> LB[Order-Taker<br/>Load Balancer]:::cloud
    LB --> T1[Truck 1<br/>App Server]:::service
    LB --> T2[Truck 2<br/>App Server]:::service
    LB --> T3[Truck 3<br/>App Server]:::service
    T1 --> DB[Recipe Book<br/>Database]:::data
    T2 --> DB
    T3 --> DB

Now here's the catch the order-taker forces on you. Suppose a customer orders at Truck 2, and the cashier there remembers "this person paid already, their burger is coming." If that customer's next request lands at Truck 3, Truck 3 has no idea who they are. The trucks must be stateless: no truck is allowed to keep private memories of customers. Any memory — who's logged in, what's in their cart — has to live somewhere every truck can reach: the shared database, or a shared cache.

This is not a minor detail; it's one of the core rules of building scalable services. Adam Wiggins wrote it down as one of the twelve factors of app design: processes should be stateless and share-nothing, with any state that needs to persist kept in a backing service (the recipe book lives outside the trucks, never inside them). (There's a whole lesson on load balancers waiting for you if you want the deeper cut — the algorithms, the health checks, the failure handling.)

Watch the requests stream through:

Interactive diagram: PacketFlow (loads in the app)

With three trucks and an order-taker, you can add a fourth, fifth, twentieth truck with almost no changes. That's the superpower of horizontal scaling: capacity becomes a number you dial up, not a wall you hit.

Video: What is a Load Balancer? — IBM Technology
The load balancer plays host at the restaurant door — greeting each request and seating it at the least-busy server so no single cook gets buried.

Stage 2: Stop re-reading the recipe book (read replicas)

The trucks multiply beautifully. But notice the diagram: all three trucks still share one recipe book — one database. And databases have a personality trait: they handle writes carefully and slowly (every write must be recorded faithfully), but they can hand out reads much more easily.

Most real services are wildly read-heavy. A social feed might serve a thousand profile views for every profile update. So the database's bottleneck is usually reads, and there's a beautiful trick for it: read replicas. Photocopy the recipe book. Put a copy at every truck. Reads go to the nearest copy; writes still go to the one true original, which then updates all the copies.

flowchart TD
    T[Trucks<br/>App Servers]:::service --> LB[Read/Write Splitter]:::cloud
    LB -->|Writes: orders, payments| P[Master Recipe Book<br/>Primary DB]:::data
    LB -->|Reads: menus, profiles| R1[Photocopy 1<br/>Read Replica]:::data
    LB -->|Reads: menus, profiles| R2[Photocopy 2<br/>Read Replica]:::data
    P -->|copies changes over| R1
    P -->|copies changes over| R2

In the food-truck world: nobody needs to walk back to the original recipe book to check the menu. The copies are fine — as long as you accept that a just-changed special might take a minute to reach every photocopy. That little delay is called replication lag, and it's a trade-off you choose on purpose: slightly stale reads in exchange for massively more read capacity.

Two honest warnings. First, writes still funnel through the one master, so if your problem is write volume, replicas don't save you. Second, that lag is real — if a customer changes their password and the replica still has the old one, they'll have a confusing minute. Some systems route a user's own reads to the master right after they write ("read your own writes") to dodge exactly this.

Video: Distributed Consensus and Data Replication strategies on the server — Gaurav Sen
Covers primary/replica split, replication lag, and read-your-own-writes with clear visual explanations.

Stage 2.5: The chalkboard by the window (caching)

Before we leave the database behind, one more trick — and it's the one with the biggest payoff on this whole list. Some questions get asked constantly: "what's today's special?", "is the truck open?" Making the cook look it up in the recipe book every time is silly when the answer barely changes. So you put a chalkboard by the window with the top ten answers written on it. The cashier reads the chalkboard; the cook is only bothered when the answer isn't there.

That's a cache: a small, blazing-fast store (usually in-memory, systems like Redis or Memcached) sitting in front of the database, holding copies of the hottest data. The usual pattern even has a folksy name — cache-aside:

  1. App server needs an answer → checks the chalkboard first.
  2. Hit: it's there. Served in microseconds, database never involved.
  3. Miss: not there. Ask the database, then write the answer on the chalkboard for next time.
sequenceDiagram
    participant Truck as App Server
    participant Board as Chalkboard (Cache)
    participant Book as Recipe Book (DB)

    Truck->>Board: What's today's special?
    Board-->>Truck: Hit! "Fish tacos" (microseconds)
    Truck->>Board: What's the soup?
    Board-->>Truck: Miss — never heard of it
    Truck->>Book: What's the soup?
    Book-->>Truck: "Clam chowder" (milliseconds)
    Truck->>Board: Write it down: soup = clam chowder

Two honest footnotes. First, chalkboards get stale: if the special changes, the board still says fish tacos until someone erases it. Caches handle this with a TTL (time to live) — every chalk note auto-erases after, say, 60 seconds. Short TTL: fresher answers, more database trips. Long TTL: fewer trips, staler answers. Another dial you tune on purpose.

Second, a chalkboard is small by design. When it's full, something gets erased to make room (eviction — usually the least-recently-used note). A cache doesn't replace your database or your replicas; it protects them, absorbing the repetitive questions so the expensive machinery underneath can breathe. There's a whole lesson on caching strategies later — consider this your appetizer.

Video: Redis Cache Explained in 300 Seconds (Beginner to Pro) #redis #cache #coding #seniordeveloper — Coding With Raj Dave
Covers cache-aside, TTL, and eviction policies with Redis in a compact format.

Stage 3: Let the corner stores sell the drinks (CDN)

By now your trucks serve food fast. But half of what customers ask for isn't cooked to order — it's the bottled drinks, the pre-made desserts, the printed menus. In web terms: images, videos, stylesheets, JavaScript bundles. Static assets. They never change per customer, so why is your precious cook fetching them from the truck?

A CDN (content delivery network) is a franchise of corner stores that stock your drinks for you. You make one delivery to the franchise; after that, customers grab drinks from the store on their corner instead of driving to your truck. Technically: copies of your static files sit on servers all over the world, and each user downloads them from the nearest one.

flowchart LR
    U[Customers]:::client --> CDN[Corner Stores<br/>CDN]:::cloud
    U --> LB[Order-Taker<br/>Load Balancer]:::cloud
    CDN -->|miss: fetch once| Origin[Truck HQ<br/>Origin Server]:::service
    LB --> S[App Servers]:::service
    S --> DB[(Database)]:::data

The wins are twofold. Your trucks stop wasting effort on work any corner store could do, and customers get their drinks faster because the corner store is closer than your truck. For media-heavy services, a CDN isn't an optimization — it's the difference between possible and impossible. (There's a full lesson on CDNs later in this course.)

One catch: when you change a drink's recipe (update an asset), the corner stores keep selling the old bottles until their stock refreshes. That's cache invalidation, and there's a famous joke that it's one of the two hard problems in computer science. The usual fix is versioned filenames — ship menu-v2.png instead of overwriting menu.png — so old stock and new stock never get confused.

Video: What is a CDN? - #Cloudbits Episode 2 — Cloudflare Developers
Cloudflare's own explainer on what CDNs are and how edge caching works.

Stage 4: Split the recipe book (sharding)

Everything scales now except the one master recipe book. Writes keep growing, and one database — even the beefiest — eventually can't keep up. The final stage of the classic journey: shard the data. Tear the recipe book into volumes and give each volume to a different truck.

Sharding means splitting your data across multiple databases by some key: customers A–M on one shard, N–Z on another; or by geography, or by a hash of the user ID. Each shard is an independent database holding a slice of the whole. Now writes scale too — adding a shard adds write capacity.

flowchart TD
    T[Trucks<br/>App Servers]:::service --> R[Shard Router]:::cloud
    R -->|users A–M| S1[Volume 1<br/>Shard]:::data
    R -->|users N–Z| S2[Volume 2<br/>Shard]:::data
    R -->|users 0–9| S3[Volume 3<br/>Shard]:::data

But — and here the "but" is doing honest work — sharding is the most expensive stage on this list. Queries that used to touch one database now might touch many ("find every customer named Lee" means asking all three volumes). Rebalancing shards as data grows is delicate surgery. And some operations, like transactions across shards, range from painful to genuinely impossible. Shard when the math says you must, not before. (There's a dedicated lesson on sharding strategies for exactly that math.)

Notice something about the whole journey: each stage attacked the current bottleneck and nothing else. You didn't shard on day one. You didn't need a CDN for your first hundred users. Scaling is sequential problem-solving, not a shopping list.

Video: System Design Part 5: Indexing & Sharding Explained (Scale to Millions) — Coding Devs
Covers sharding strategies, shard keys, hot shards, and cross-shard query costs for write scaling.

The napkin math

Senior engineers do capacity planning on the back of a napkin, and you can too. The Google SRE book treats this as a core discipline: translate users into requests per second, then into machines. Here's the whole technique in one worked example.

Step 1: Users → requests per second. Your analytics say the service handles 1 million requests per day. A day has 86,400 seconds, so:

That "10× peak" rule of thumb is deliberately pessimistic. Underestimate the peak and your launch-day story becomes a postmortem.

Step 2: Requests per second → servers. Suppose one app server comfortably handles ~500 requests/second. Your peak of ~120 req/s fits on a single server — so why did we build a fleet of three? Because one server is a single point of failure, and because the peak estimate is a guess. You run three for redundancy and headroom, not because the math demands three. Capacity planning is half arithmetic, half humility.

Step 3: Data → storage. Storage math is even simpler — multiply the thing by the count:

That last line is the real lesson of this section. Back-of-the-envelope math won't give you the exact answer, but it reliably tells you which stage of the journey you're in — and which stages you can stop worrying about for now.

Interactive diagram: NapkinMathPlayground (loads in the app)

Video: SREcon19 Europe/Middle East/Africa - Advanced Napkin Math: Estimating System... — USENIX
SREcon talk on estimating system performance from first principles, matching the napkin-math discipline.

The full journey, one diagram

Here's everywhere the food truck has been, in the order it got there:

flowchart TD
    U[Customers]:::client --> CDN[CDN<br/>static assets]:::cloud
    U --> LB[Load Balancer]:::cloud
    LB --> A1[App Server]:::service
    LB --> A2[App Server]:::service
    LB --> A3[App Server]:::service
    A1 --> R[Read/Write Splitter]:::cloud
    A2 --> R
    A3 --> R
    R -->|reads| RR[Read Replicas]:::data
    R -->|writes| SH[Sharded Primaries]:::data

Stateless trucks behind an order-taker. Reads fanned out to photocopies. Writes spread across volumes. Static assets handled by the franchise. Every stage earned its place by solving the bottleneck in front of it.

Now walk the whole journey yourself — scroll, and watch each stage light up in the order the truck actually grew:

Interactive diagram: ScrollyDiagram (loads in the app)

Video: System Design Explained: How Apps Scale to 1 Billion Users — TheBinaryBreakdown
Rebuilds the whole journey from one server to global architecture: load balancers, caching, replicas, sharding, CDN.

Takeaways

  1. Scale vertically first, horizontally second. A bigger truck is simple but has a hard ceiling and a single point of failure; a fleet of trucks scales further and survives breakdowns.
  2. Horizontal scaling demands stateless servers. A load balancer can send any request to any server, so no server may keep private memories — shared state lives in the database or cache, never in the app process.
  3. Read replicas multiply read capacity by photocopying the database; writes still flow through the primary. Accept a little replication lag as the price.
  4. A CDN offloads static assets to servers near your users — faster for them, cheaper for you. Version your filenames to dodge cache invalidation headaches.
  5. Shard only when the math says so. Splitting the database scales writes, but cross-shard queries and rebalancing are genuinely painful. The napkin math tells you when you've arrived.
  6. Do capacity planning on a napkin: daily requests ÷ 86,400 = average req/s, assume ~10× for peak, then size the fleet with room to spare. Half arithmetic, half humility.

Check your understanding

  1. Your food truck's line is out the door. What's the key difference between scaling vertically and scaling horizontally?

    • Vertical means hiring more cooks; horizontal means opening longer hours
    • They're two names for the same thing
    • Vertical means a bigger, more powerful truck; horizontal means more trucks side by side
    • Vertical means adding more trucks; horizontal means getting a bigger grill
  2. Why must the app servers (trucks) behind a load balancer be stateless?

    • Stateless servers use less electricity
    • The load balancer can route any customer's next request to a different server, so a server can't rely on remembering that customer
    • Databases refuse connections from servers that keep local state
    • Stateless servers never need to be restarted
  3. Your service handles 1M requests/day. What are the average and (10×) peak requests per second?

    • About 12 avg, about 120 peak
    • About 60 avg, about 600 peak
    • About 120 avg, about 1,200 peak
    • About 1,000 avg, about 10,000 peak
  4. Read replicas help you scale, but they do NOT help with which problem?

    • Serving reads from a database copy close to the app servers
    • Reducing load on the primary database from reads
    • Serving a huge volume of read queries
    • Handling a huge volume of writes — writes still funnel through the single primary

Go deeper

Want to keep pulling this thread? These talks and tutorials go further than we did here:

Sources & further reading