Your country

Tools that support it use your country for local currency, number formats, units and paper size. Your choice is saved only in this browser.

Type a name or a two-letter code. Use the up and down arrow keys to move through the countries, Enter to choose one and Escape to close.

System Design (High-Level Design) Module 1 – Start here: how system design works

Growing a service from one server to many regions

How a service grows from one server to many regions, what breaks first at each stage, which component fixes it, and why adding parts too early costs you.

  • Beginner
  • 30 minutes
  • Examples run with Python 3.14.8, Pyodide 314.0.7, Node.js 24.21.0 and quickjs 0.32.0
  • By MySmartCoPilot

What you will learn

  • Trace how a service evolves as traffic grows by orders of magnitude
  • Explain what breaks first at each stage and which component fixes it
  • Avoid adding distributed components before the load needs them

Before you start

On this page

Almost every large service started small. The design that fits a hundred users would collapse under ten million, and the design for ten million would bury a small team serving a hundred. This lesson follows one service through the stages services usually pass on the way up, says what runs out first at each stage and which component fixes it, and then turns the picture around: every component you add has a price, and paying it too early is a design mistake.

Stage 1: one server

Stage 1: users reach one server that runs the app, the database and the static files together.UsersOne serverAppDatabaseImages and scriptson its diskSQL, same machinereadsHTTPS: pages, files, API

Stage 1: everything on one server

Text description of the diagram

Users send requests for pages, files and the API over HTTPS to one server. Inside that server: - the app answers every request; - it talks SQL to a database on the same machine; - it reads images and scripts from the server's own disk.

Everything runs on one machine: the app, the database and the files. It is cheap, simple to deploy and easy to debug, because there is only one place to look. For a prototype, an internal tool or a young product it is often the right design, and for longer than people expect. What breaks first is that everything competes for the same CPU, memory and disk, and that the machine is a single point of failure: a crash, a bad deploy or a full disk takes the whole service down, and a backup is the only protection for the data.

Each stage fixes the bottleneck of the one before

Each stage adds one thing, because of one bottleneck, and has a price:

  1. A separate database server, because the app and the database compete for one machine. It costs a network hop on every query and a second machine to run.
  2. A load balancer and several stateless app servers, because one app server’s CPU is full and every deploy takes the site down. State must leave the servers, the load balancer needs health checks, and it is one more component.
  3. A CDN for static files, because the servers spend their bandwidth sending the same images and scripts. Old files stay cached far away, so file names carry versions, and there is one more provider.
  4. A cache in front of the database, because the database answers the same reads again and again. It brings stale data, invalidation, and a cold cache after a restart.
  5. Read replicas, because reads still outgrow one database. Replication lag means a user may not see their own change at once.
  6. A queue and workers for slow work, because emails, thumbnails and reports make requests slow, or make them fail when a provider is slow. That work now finishes later, may be retried or done twice, and leaves a backlog to watch.
  7. Sharding, because writes or the size of the data outgrow one primary database. Joins and transactions across shards stop being easy, and data has to move as shards fill.
  8. More regions, because distant users wait on every request and one region’s outage takes everything down. Data must be copied between regions, writes conflict or slow down, and running the service gets much harder.

Each stage has a source you can check. Shared caches in front of a service exist to cut response time and bandwidth by reusing stored responses (the HTTP caching standard, RFC 9111). With the cache-aside pattern, the app loads data into the cache on demand and has to decide how long an entry may live: too short and the database still does the work, too long and readers see stale data. A PostgreSQL standby answers read-only queries, and because changes reach it after a delay, the same query on the primary and the standby can return different answers. A queue between a task and a slow service lets the service work at its own pace and be sized for the average load instead of the peak. And a single database server eventually runs out of storage, computing power or network bandwidth, which is when the data is divided into shards that each hold part of it.

By stage 8 the same service, still in one region, looks like this:

One region grown to many machines: CDN, load balancer, stateless app servers, cache, database shards, a queue and workers.UsersCDNimages, scripts, stylesLoad balancerApp serversmany, identical, statelessCachehot readsDatabase shardsa primary and replicas eachQueueslow jobsWorkersemails, thumbnails, reportsObject storageuploads, static filesstatic filespages and API, HTTPSfetch on a missany server, any request1. read2. on a miss; every writehand off slow jobsdeliver jobsstore results

Later: the same service spread over many machines in one region

Text description of the diagram

Users reach the service in two ways: - static files (images, scripts and styles) come from a CDN, which fetches anything it does not hold from object storage; - pages and API calls go over HTTPS to a load balancer, which may send any request to any of many identical, stateless app servers.

The app servers first read from a cache of hot data (step 1). On a cache miss, and for every write, they go to the database, which is split into shards, each with a primary and replicas (step 2).

Slow jobs, such as sending emails, making thumbnails and building reports, are handed to a queue, which delivers them to workers; the workers store their results in object storage.

Stateless servers make scaling out easy

Stage 3 is the one that unlocks most of the rest. A stateless app server keeps nothing between requests that another server would need: sessions live in a shared store or a signed token (the session trade-off of an earlier lesson), uploads go to object storage, and anything that must last goes to a database. Then any server can answer any request, so the load balancer can send traffic anywhere, a failed server can simply be replaced, and new versions can be deployed one server at a time while the others keep serving.

The Twelve-Factor App states the rule bluntly: its processes are stateless and share nothing, anything that must persist lives in a backing service such as a database, and sticky sessions are not allowed. It also explains why the rule pays: one machine can only grow so large, but share-nothing processes can be multiplied. AWS’s reliability guidance makes the same point from the failure side: replace one large resource with several small ones, so that losing one has a small effect.

Watch it grow: a toy model

This program raises the traffic of two different services tenfold at a time. At each step it works out how busy each resource is, fixes the busiest one that is above 70% (the rest is headroom for peaks and for losing a machine), and repeats until everything fits. Every capacity at the top is an assumption, not a measurement.

What runs out first as traffic grows Python · growth_sim.py
"""A toy model of a service that grows: raise the traffic tenfold at a time, find the resource that is busiest
past the headroom limit, apply the usual fix for it, and repeat until everything fits again.
Every capacity below is an assumption for the sketch, not a measurement: change them and run it again."""

import math


def count(n, word):
    """'1 write', '2 writes', '0.05 writes'."""
    return f"{n} {word}" if n == 1 else f"{n} {word}s"

HEADROOM = 0.7  # keep every resource under 70% busy, to absorb peaks and the loss of one machine
APP_SERVER = 500  # requests a second one app server handles
DB_READS = 3000  # simple reads a second one database server answers
DB_WRITES = 1000  # writes a second one database server takes
NETWORK = 1250  # MB a second one load balancer can send (10 Gbit/s)
SHARED = 0.5  # app and database on one machine: each gets about half of it
CACHE_HITS = 0.9  # share of reads a cache answers, once there is one

WORKLOADS = {
    "A read-heavy site": {"reads": 2, "writes": 0.05, "kb": 100, "static_kb": 80},
    "A write-heavy tracker": {"reads": 0.2, "writes": 1, "kb": 2, "static_kb": 0},
}


def busy(rate, w, s):
    """How busy each resource is (1.0 = full) at `rate` requests a second."""
    share = 1 if s["own_db"] else SHARED
    reads = rate * w["reads"] * (1 - CACHE_HITS if s["cache"] else 1)
    sent_kb = w["kb"] - (w["static_kb"] if s["cdn"] else 0)
    return {
        "app servers": rate / (APP_SERVER * share * s["app"]),
        "database reads": reads / (DB_READS * share * (s["shards"] + s["replicas"])),
        "database writes": rate * w["writes"] / (DB_WRITES * share * s["shards"]),
        "network": rate * sent_kb / 1000 / (NETWORK * s["balancers"]),
    }


def fix(resource, rate, w, s):
    """Change the design so that `resource` has room again; say what was done."""
    if resource != "network" and not s["own_db"]:
        s["own_db"] = True
        return "move the database to a server of its own"
    if resource == "app servers":
        s["app"] = math.ceil(rate / (APP_SERVER * HEADROOM))
        return f"run {s['app']} stateless app servers behind a load balancer"
    if resource == "database reads" and not s["cache"]:
        s["cache"] = True
        return "put a cache in front of the database"
    if resource == "database reads":
        reads = rate * w["reads"] * (1 - CACHE_HITS)
        s["replicas"] = math.ceil(reads / (DB_READS * HEADROOM)) - s["shards"]
        return f"add read replicas: {s['replicas']} in all"
    if resource == "database writes":
        before = s["shards"]
        s["shards"] = math.ceil(rate * w["writes"] / (DB_WRITES * HEADROOM))
        s["replicas"] = max(0, s["replicas"] - (s["shards"] - before))  # each new shard serves reads too
        return f"split the data by key into {s['shards']} shards"
    if not s["cdn"]:
        s["cdn"] = True
        return "serve images, scripts and styles from a CDN"
    s["balancers"] = math.ceil(rate * (w["kb"] - w["static_kb"]) / 1000 / (NETWORK * HEADROOM))
    return f"spread traffic over {s['balancers']} load balancers"


for name, w in WORKLOADS.items():
    print(f"{name}: per request {count(w['reads'], 'read')}, {count(w['writes'], 'write')} and {w['kb']} KB sent")
    s = {"own_db": False, "app": 1, "cache": False, "replicas": 0, "shards": 1, "cdn": False, "balancers": 1}
    for rate in (10, 100, 1_000, 10_000, 100_000):
        steps = []
        while True:
            load = busy(rate, w, s)
            resource = max(load, key=load.get)
            if load[resource] <= HEADROOM:
                break
            steps.append(f"{resource} at {load[resource]:.0%}: {fix(resource, rate, w, s)}")
        busiest = max(load, key=load.get)
        print(f"  {rate:,} requests a second, busiest {busiest} at {load[busiest]:.0%}")
        for step in steps:
            print(f"    {step}")
    parts = [count(s["app"], "app server"), count(s["shards"], "database shard"), count(s["replicas"], "read replica")]
    parts += [label for label, used in (("a cache", s["cache"]), ("a CDN", s["cdn"])) if used]
    print(f"  In the end: {', '.join(parts)} and {count(s['balancers'], 'load balancer')}")
    print()

Output

A read-heavy site: per request 2 reads, 0.05 writes and 100 KB sent
  10 requests a second, busiest app servers at 4%
  100 requests a second, busiest app servers at 40%
  1,000 requests a second, busiest app servers at 67%
    app servers at 400%: move the database to a server of its own
    app servers at 200%: run 3 stateless app servers behind a load balancer
  10,000 requests a second, busiest app servers at 69%
    app servers at 667%: run 29 stateless app servers behind a load balancer
    database reads at 667%: put a cache in front of the database
    network at 80%: serve images, scripts and styles from a CDN
  100,000 requests a second, busiest app servers at 70%
    app servers at 690%: run 286 stateless app servers behind a load balancer
    database reads at 667%: add read replicas: 9 in all
    database writes at 500%: split the data by key into 8 shards
    network at 160%: spread traffic over 3 load balancers
  In the end: 286 app servers, 8 database shards, 2 read replicas, a cache, a CDN and 3 load balancers

A write-heavy tracker: per request 0.2 reads, 1 write and 2 KB sent
  10 requests a second, busiest app servers at 4%
  100 requests a second, busiest app servers at 40%
  1,000 requests a second, busiest app servers at 67%
    app servers at 400%: move the database to a server of its own
    app servers at 200%: run 3 stateless app servers behind a load balancer
    database writes at 100%: split the data by key into 2 shards
  10,000 requests a second, busiest app servers at 69%
    app servers at 667%: run 29 stateless app servers behind a load balancer
    database writes at 500%: split the data by key into 15 shards
  100,000 requests a second, busiest app servers at 70%
    app servers at 690%: run 286 stateless app servers behind a load balancer
    database writes at 667%: split the data by key into 143 shards
  In the end: 286 app servers, 143 database shards, 0 read replicas and 1 load balancer

Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 growth_sim.py

Read the first service, a read-heavy site, from the top:

  • At 10 and 100 requests a second nothing is added. One server is enough, and anything more would be cost and work for nothing.
  • At 1,000, the shared machine runs out first, so the database moves to its own server; then the app needs three servers behind a load balancer.
  • At 10,000, the database would be answering the same reads 20,000 times a second, so a cache goes in front of it, and the network fills with images and scripts, so they move to a CDN.
  • At 100,000, even the reads that miss the cache outgrow one database, writes outgrow one primary, and the data is split into shards with read replicas.

The second service, a tracker that mostly writes small records, grows differently. It never needs a cache or a CDN, but the database’s writes run out at 1,000 requests a second, and by 100,000 the model wants 143 shards. That number is the real lesson: when the only fix you know is “more shards”, a write-heavy workload needs a different design, such as batching writes, putting them through a log, or a store built for heavy writes, which later modules cover. The order of the stages is not a law. It follows the workload, which is why the estimates of a design round come before the diagram.

The model is a sketch for reasoning, with one fix per resource; it knows nothing about latency, failures or cost. Change the capacities or the workloads at the top and run it again to see how the order moves.

Do not add parts before the load needs them

Every stage on that list adds failure modes, work for whoever is on call, and money. Azure’s guidance on reliability puts it plainly: each component added to make a system resilient also makes it more complex to run. In concrete terms:

  • A cache brings invalidation bugs, stale reads, and a stampede on the database when it restarts empty.
  • A read replica brings users who do not see the change they just made.
  • A queue brings work that happens later, twice, or not at all, and a backlog someone must watch.
  • Shards take away joins and transactions across shards and bring the work of moving data when one fills up.
  • More regions bring conflicting or slower writes and double the infrastructure.

So add a component when a measured or estimated bottleneck needs it, and write down the trigger with it: “add a read replica when the database’s CPU stays above 60% at the daily peak”. In a design round, a design for a few thousand users with a message log, a cache cluster, sixteen shards and three regions is not a strong answer. A strong answer says what the design looks like now, and which component it would add at which number.

A first-hand example

This site skips most of the stages. Every page is built ahead of time and served as static files from a content delivery network, most tools run inside the visitor’s own browser, and only a few small server functions with one small SQL database handle what needs shared state, such as contact forms. A read-heavy service whose work can run on the user’s device can go a long way without app servers of its own.

Many regions

The last stage deserves its own warning. Going to several regions makes sense when users far away pay for the distance on every request (no hardware makes a signal cross an ocean faster), when the service must survive the loss of a whole region, or when rules say where data must stay. The hard part is always the data: which region accepts writes, and how they reach the others. The write-path trade-off of the quality-attributes lesson applies with large numbers: synchronous copies between regions make every write wait for an ocean crossing, asynchronous ones can lose the last writes or let two regions accept conflicting ones. The two common shapes are active-passive, where one region takes the writes while another stands ready, and active-active, where every region serves users, which is faster for them and much harder to keep correct. A later module of this track covers both in depth.

HTTP Status Codes Reference Look up 429 and 503, the answers an overloaded service should give instead of timing out.

Key takeaways

  • Services grow in stages: one server, a separate database, stateless app servers behind a load balancer, a CDN, a cache, read replicas, queues for slow work, shards, then more regions.
  • Each stage removes the bottleneck of the one before and adds failure modes, operational work and cost.
  • Stateless app servers are what make scaling out, replacing servers and rolling deploys easy.
  • The order depends on the workload: read-heavy services lean on caches and CDNs, write-heavy ones on the write path.
  • Add a component when a measured or estimated bottleneck needs it, and write down the number that triggers it.

Exercise

Exercise · Easy · Python

Find the resource that fills up first

As traffic grows, the resource that fills up first is the one with the least capacity for each request. Write bottleneck(per_request, capacity) in bottleneck.py to find it.

- per_request says how much of each resource one request uses, such as {"app_cpu_ms": 4, "db_reads": 2}: 4 milliseconds of app server CPU and 2 database reads. - capacity says how much of each resource there is per second, such as {"app_cpu_ms": 8000, "db_reads": 3000}: 8 CPU cores give 8,000 milliseconds of CPU a second.

Return a pair: the name of the resource that is full at the lowest request rate, and that rate in requests a second, rounded to one decimal. Here the CPU is full at 8,000 / 4 = 2,000 requests a second and the database at 3,000 / 2 = 1,500, so:

bottleneck({"app_cpu_ms": 4, "db_reads": 2}, {"app_cpu_ms": 8000, "db_reads": 3000})  # ('db_reads', 1500.0)

A resource that a request does not use at all (0) never fills up, so leave it out. When two resources fill up at the same rate, return the one whose name comes first in alphabetical order. If a request uses no resource at all, raise ValueError. Every resource in per_request has an entry in capacity.

The sample tests import bottleneck from bottleneck.py and run in your browser.

Starter code · bottleneck.py

def bottleneck(per_request, capacity):
    """The resource that is full at the lowest request rate, and that rate rounded to 1 decimal."""
    # Replace this line with your code.
    return ("", 0.0)
The sample tests · test_bottleneck.py
from bottleneck import bottleneck


def raises_value_error(call):
    try:
        call()
    except ValueError:
        return True
    return False


def test_the_database_fills_first():
    """finds the resource with the least capacity for each request"""
    assert bottleneck({"app_cpu_ms": 4, "db_reads": 2}, {"app_cpu_ms": 8000, "db_reads": 3000}) == ("db_reads", 1500.0)


def test_the_cpu_fills_first():
    """works when the first resource is the bottleneck"""
    assert bottleneck({"app_cpu_ms": 5, "db_reads": 1}, {"app_cpu_ms": 4000, "db_reads": 3000}) == ("app_cpu_ms", 800.0)


def test_unused_resources_never_fill():
    """leaves out a resource that a request does not use"""
    assert bottleneck({"cpu_ms": 0, "network_kb": 50}, {"cpu_ms": 1, "network_kb": 125000}) == ("network_kb", 2500.0)


def test_ties_and_rounding():
    """settles a tie alphabetically and rounds the rate to one decimal"""
    assert bottleneck({"writes": 2, "reads": 1}, {"writes": 200, "reads": 100}) == ("reads", 100.0)
    assert bottleneck({"x": 3}, {"x": 1000}) == ("x", 333.3)


def test_no_resource_used():
    """raises ValueError when a request uses nothing"""
    assert raises_value_error(lambda: bottleneck({"cpu_ms": 0}, {"cpu_ms": 1000}))
A hint

For each resource that a request uses (used > 0), the rate at which it is full is capacity[name] / used. Keep those rates in a dictionary, raise ValueError if it is empty, and pick the name with the smallest rate. Going through the names in sorted order and keeping the first smallest one settles ties alphabetically.

The sample tests run on this device, in your browser (Pyodide): nothing is sent to mysmartcopilot.com. The first run downloads Python (about 13.5 MB), which is kept for the next runs. A check in your browser is feedback for you, not proof that the code is right for every input.

Check yourself

7 questions about this lesson. Every answer and why it is right is on the page, behind “Show the answer”. Your score stays in this browser.

  1. Question 1 of 7 Your one server's network link is saturated at the evening peak, while its CPU is at 30%. Most bytes it sends are the same images and scripts. What do you add next?

    Choose one answer.

    Show the answer to question 1

    Answer: A CDN that serves the static files from caches near the users

    The bottleneck is bandwidth spent on identical static files, which shared caches exist to absorb. More app servers or replicas would not reduce the bytes this link sends.

  2. Question 2 of 7 The app servers run at 95% CPU at the peak while the database sits at 20%. The servers keep sessions in a shared store. What do you add?

    Choose one answer.

    Show the answer to question 2

    Answer: More app servers behind the load balancer

    The app tier is the bottleneck and the servers are already stateless, so adding servers is the direct fix. The database has plenty of room, so caching or sharding it would not help yet.

  3. Question 3 of 7 The database spends most of its time answering the same few product-page queries, and the products change only a few times a day. What do you add?

    Choose one answer.

    Show the answer to question 3

    Answer: A cache in front of the database for those reads

    Repeated reads of data that rarely changes are exactly what a cache absorbs, and a slightly old product page is acceptable here. Sharding is for writes or data size, not for the same reads repeated.

  4. Question 4 of 7 Checkout waits three seconds while a confirmation email is sent, and checkouts fail whenever the email provider is slow. What do you add?

    Choose one answer.

    Show the answer to question 4

    Answer: A queue and a worker that sends the email after the checkout has been answered

    The email does not need to happen inside the request. Handing it to a queue lets checkout answer at once and lets the worker retry while the provider is slow; more servers would only wait for the provider in parallel.

  5. Question 5 of 7 The primary database is at 90% from writes alone; reads already go to replicas and a cache. What is the next step?

    Choose one answer.

    Show the answer to question 5

    Answer: Split the data by key into shards, so that each primary takes part of the writes

    Replicas and caches take load off reads, not writes: every write still goes through the one primary. Sharding divides the writes between several primaries (batching writes is another way to buy room).

  6. Question 6 of 7 Users on another continent wait an extra 250 ms on every request, and the business now needs to survive the loss of a whole region. What do you add?

    Choose one answer.

    Show the answer to question 6

    Answer: A second region, with a plan for which region accepts writes and how data is copied between them

    Distance and regional outages cannot be fixed inside one region. A second region brings servers closer to those users and survives the loss of the first, and the hard part is deciding how writes and data move between them.

  7. Question 7 of 7 An internal tool for 50 colleagues runs on one server at 5% CPU. A teammate proposes a message log, a cache cluster and sharding "to be ready for growth". What is the best answer?

    Choose one answer.

    Show the answer to question 7

    Answer: Keep one server, add backups and monitoring, and write down the numbers at which each component would be added

    None of those components removes a bottleneck the tool has, and each adds failure modes, on-call work and cost. Protecting the data with backups and recording the triggers for growth is the proportionate design.

References

Related tools

Report a problem with this lesson

Quick answers and tool search

Type to search tools or to get a quick answer, for example 18% of 2500. Use the up and down arrow keys to move through the results, Enter to choose, and Escape to close.