Module 1 · Thinking in scale

The journey: 1 user to 1 million

Foundations 16 min read What "scale" means, in numbers you can check

Meet Linkly, a URL shortener. Its founder writes the first version over a weekend: one small server, one database, one happy user. This course follows Linkly all the way to a million users. At every stage something specific breaks, we measure it, and we add exactly the piece of architecture that fixes it. By the end you will know not just what a load balancer, a cache or a shard is, but when each one earns its place and what it costs.

🍜

From a food stall to a restaurant chain

A food stall does not start with a central kitchen, a delivery fleet and a franchise manual. It starts with one cook. When the queue gets long, it hires a second cook (more servers). When the same dish is ordered all day, it pre-cooks it (a cache). When people come from the other side of town, it opens a branch (a CDN, then a region). Each change answers a problem that actually happened. Build the chain on day one and you go bankrupt before the first customer arrives.

1 user 100 1K 10K–100K 1M beyond stage 0 stage 1 stage 2 stage 3 stage 4 stage 5 app + databaseone machine app server database load balancer app × 2 database CDN load balancer app × N cache primary + replicas queue + workers CDN gateway + limits autoscaled apps cache cluster sharded database event stream 2+ regionsgeo routing stage 4, copiedin each region what forces the change: one crash loseseverything peak traffic, anddeploys = downtime database reads,slow far away writes, data size,abuse, failures a region outage,global latency Bold = added at this stage. Each stage keeps everything before it; nothing is added before something forces it.

1. What "scale" actually means

"Can it scale?" hides three different questions. A system can be fine on one axis and failing on another, so always ask which one you mean.

AxisMeasured inWhat breaks firstTypical fixes
Loadrequests per second, concurrent connectionsCPU of the app tier, database connections, one hot keymore servers, caches, queues, rate limits
Datagigabytes, rows, write ratedisk, backups taking hours, queries that scan, index size vs RAMindexes, archiving, partitioning, sharding, object storage
Organisationengineers, teams, deploys per daymerge conflicts, slow releases, "who owns this?"modules, clear ownership, sometimes services

There is also a quality you get for free at stage 0 and lose as you grow: simplicity. Every box you add is something that can fail, must be monitored, and must be understood by the next engineer.

2. From users to requests

"A million users" sounds enormous. What does it mean for the servers? We need a few assumptions about how people use Linkly. Every one of them is a guess you would check against real analytics later. The point is to write them down so they can be argued with.

from simkit import table, human, human_bytes

STAGES = [1, 100, 1_000, 10_000, 100_000, 1_000_000]   # registered users

DAU_SHARE = 0.20          # 20% of users are active on a given day
LINKS_PER_DAU = 2         # each active user creates two short links a day (writes)
CLICKS_PER_DAU = 100      # their audiences click 100 times a day in total (reads: redirects)
PEAK_FACTOR = 5           # the busiest hour runs at 5x the daily average
LINK_BYTES = 500          # code + long URL + owner + timestamps + index overhead
CLICK_EVENT_BYTES = 100   # one analytics record per click: time, code, country, referrer

rows = []
for users in STAGES:
    dau = users * DAU_SHARE
    writes, reads = dau * LINKS_PER_DAU, dau * CLICKS_PER_DAU
    avg_rps = (reads + writes) / 86_400                  # seconds in a day
    rows.append([human(users, 0), human(dau), human(writes), human(reads), f"{avg_rps:.2f}",
                 f"{avg_rps * PEAK_FACTOR:.1f}", human_bytes(writes * 365 * LINK_BYTES),
                 human_bytes(reads * 365 * CLICK_EVENT_BYTES)])
table(rows, ["users", "DAU", "links/day", "clicks/day", "avg req/s", "peak req/s", "links/yr", "click log/yr"])
users DAU links/day clicks/day avg req/s peak req/s links/yr click log/yr ----- ------ --------- ---------- --------- ---------- -------- ------------ 1 0.2 0.4 20 0.00 0.0 73.0 KB 730.0 KB 100 20 40 2.0K 0.02 0.1 7.3 MB 73.0 MB 1K 200 400 20.0K 0.24 1.2 73.0 MB 730.0 MB 10K 2.0K 4.0K 200.0K 2.36 11.8 730.0 MB 7.3 GB 100K 20.0K 40.0K 2.0M 23.61 118.1 7.3 GB 73.0 GB 1M 200.0K 400.0K 20.0M 236.11 1180.6 73.0 GB 730.0 GB

The surprise

A million users is about 240 requests per second on average and roughly 1,200 at peak. A single well-tuned server can handle that. The reasons to go beyond one server are elsewhere.

Reads dominate

With 50 clicks for every new link, this is a read-heavy system. That points at caches and replicas long before sharding.

Data grows quietly

The links themselves stay small (73 GB a year). The click log does not: 730 GB a year, and it never shrinks. Analytics becomes the data problem.

3. Averages lie: spikes, hot keys and availability

If average load is small, why do real systems at this size need load balancers, caches and replicas? Because traffic is not average. Here is one link going viral, and what three years of the click log look like:

from simkit import human, human_bytes

# A celebrity posts a Linkly link to 30M followers; 4% click, almost all in the first ten minutes.
followers, click_rate, window_s = 30_000_000, 0.04, 600
viral_rps = followers * click_rate / window_s
normal_peak = 1_180                                   # from the table above, at 1M users
print(f"viral link: {human(followers * click_rate)} clicks in 10 minutes = {viral_rps:,.0f} req/s "
      f"on ONE key ({viral_rps / normal_peak:.1f}x the normal peak for the whole site)")

# The click log keeps everything unless someone decides otherwise.
per_day = 20_000_000 * 100                            # clicks/day at 1M users x bytes per event
for year in (1, 2, 3):
    print(f"click log after {year} year{'s' if year > 1 else ''}: {human_bytes(per_day * 365 * year)}")

# Availability: what a single server's downtime means in a year
for nines, label in [(0.99, "99%"), (0.999, "99.9%"), (0.9999, "99.99%")]:
    minutes = (1 - nines) * 365 * 24 * 60
    print(f"{label:>7} available = {minutes:,.0f} minutes of downtime a year")
viral link: 1.2M clicks in 10 minutes = 2,000 req/s on ONE key (1.7x the normal peak for the whole site) click log after 1 year: 730.0 GB click log after 2 years: 1.5 TB click log after 3 years: 2.2 TB 99% available = 5,256 minutes of downtime a year 99.9% available = 526 minutes of downtime a year 99.99% available = 53 minutes of downtime a year
The real reasons architecture grows

Spikes (the viral link is more than the rest of the site put together, and it all hits one row), data volume (terabytes of history), availability (one server means every reboot and every deploy is an outage, and 99.9% still allows almost nine hours of downtime a year), and latency for distant users. Average requests per second is the least interesting number on the page.

4. The stages, and where this course covers them

StageUsersWhat hurtsWhat we addLessons
01–100nothing yet; speed of building matters mostone server, one process, backups04
1100–1Kapp and database fight for memory; one crash loses botha separate, managed database; indexes05
21K–10Kdeploys cause downtime; peaks overload one serverstateless apps behind a load balancer06, 07
310K–100Kdatabase reads saturate; slow pages abroad; slow requests doing heavy workcache, CDN, read replicas, queues08–11
4100K–1Mwrites and data outgrow one database; abuse; partial failures cascadesharding, rate limits, timeouts and breakers, SLOs12–18
51M+fan-out, region outages, global users, the cloud billfeeds, multiple regions, security at the edge, capacity planning19–22

5. Decide what is cheap now and expensive later

"Don't build for a million users on day one" is not "don't think". Some decisions are two-way doors (easy to change later) and some are one-way doors. Spend your early care on the one-way doors:

DecisionCheap on day onePainful at a million
Identifiersrandom or time-ordered IDs (UUIDv7, Snowflake-style)auto-increment IDs leak volume, collide across shards (lesson 12, 21)
Statesessions and uploads outside the app processrewriting login and storage while adding servers (lesson 06)
Timestore UTC everywheremigrating billions of rows of local times
Configurationenvironment variables, not hard-coded hostsevery new environment needs a code change
Data modela clear owner for each table; no cross-module joins by habituntangling a shared database to split anything (lesson 17)
Observabilityrequest IDs and structured logs from the startdebugging production blind (lesson 18)

6. How every lesson works

Symptom → measurement → fix → new bottleneck

We never add a component because diagrams usually have one. Each lesson starts from something that hurts, measures it with runnable code, applies the fix, and names the next bottleneck.

Simulations you can run

Networks of hundreds of servers don't fit in a lesson, so many examples are small seeded simulations (see simkit.py in this course folder). Where a real server fits on a laptop, the course runs a real one and measures it.

Recap

  • Scale has three axes: load, data and organisation. Name the one you mean.
  • Turn users into numbers with written-down assumptions: DAU, actions per user, read/write ratio, peak factor, bytes per record.
  • Averages are misleading: spikes, hot keys, data growth and availability drive most architecture.
  • Add components when something forces you to, but get the one-way doors (IDs, state, time, config) right early.

Checkpoint

1 · Linkly has 1M users and about 1,200 peak requests per second. Why might it still need several servers?
Average load is modest. One server is still a single point of failure, and spikes can be several times the normal peak.
2 · Which early decision is the hardest to change later?
Server counts and cache sizes are configuration. IDs end up in URLs, other systems and every foreign key, and sequential IDs collide across shards.
3 · In Linkly's numbers, which grows fastest in bytes?
Each click writes an event, and there are 50 clicks per new link. Append-only event data is usually what outgrows a single database first.