Foundations of System Design

Build the core knowledge required before tackling any system design interview question.

Last generated

Lesson 1 of 18 available15 practice questions

SPACED REPETITION Β· 15 practice questions

Make this lesson stick.

Try 3 questions now. No account needed. Sample answers aren't saved.

Watch the Lesson Video21 min
The whole lesson, narrated: the habit behind system design, what the interview grades, the life of a request, whiteboard arithmetic, and failure

One box, three million visitors

Here is a design that has launched thousands of products:

 browser ──► [ one server: web app + Postgres + images on its disk ]

It is a good design. A recipe site doing 30 page views per second at its busiest barely wakes this machine up. Adding load balancers and caches now would cost money and attention and buy no capacity you need. (What it does need from day one is backups; one box is also one failure, which comes up later.)

Then a TV cooking show mentions one of your recipes. Three million people open that page within ten minutes. That is 3,000,000 Γ· 600 s = 5,000 page views per second, about 170 times your normal peak.

To see what happens, you need three facts about a page view and three about the machine. They are assumptions. In an interview you state them out loud, and the interviewer corrects any they dislike.

Per page view Assumption The server Assumption
HTML, built by the app 100 KB Network link 10 Gbps
Images, CSS, JS 1.9 MB (2 MB in total) App CPU renders 1,500 pages/s
Database queries 4 Postgres, on the same box 8,000 simple queries/s

Predict first: which resource runs out first, and by how much? Most people guess the database.

Check your answer

Divide demand by capacity for each resource. The largest ratio breaks first.

Resource Demand at 5,000 views/s Capacity Load
Network out 5,000 Γ— 2 MB Γ— 8 bits = 80 Gbps 10 Gbps 8Γ—
App CPU 5,000 renders/s 1,500/s 3.3Γ—
Database 5,000 Γ— 4 = 20,000 queries/s 8,000/s 2.5Γ—

The network link saturates first, at eight times its capacity. Everything else is over capacity too, so a heroic database fix alone would have changed nothing a visitor could notice.

The move: stop repeating the same work

The reflex is to buy capacity: a bigger box, then eight boxes. But look at what the server is actually doing. It builds the same page three million times, runs the same four queries three million times and sends the same 1.9 MB of images, CSS and JS three million times. After the first visitor, nobody asked for anything new.

So stop repeating the work. Put a CDN (a content delivery network: caching servers in many cities) in front of the site. Let it keep the images, CSS and JS, and let it keep the HTML of public pages for 60 seconds.

  • The static files, 95% of the bytes, now come from edge servers near each reader.
  • Each edge location asks your server for the recipe about once a minute. Even with 500 edge locations, that is about 8 requests and 33 database queries per second.

One move fixed the network, the CPU and the database, because all three loads came from the same repeated work.

What it costs. If you fix a typo in the recipe, readers can see the old version for up to a minute unless you purge the CDN. Anything personal ("Hi, Sam") cannot be part of the shared cached HTML. And the CDN bills per gigabyte. For a recipe, a minute of staleness is fine. That pair, the move and the price you accept, is the basic unit of system design.

πŸ’‘ The whole course in one sentence: find the resource that runs out first, pick the cheapest move that relieves it, and say out loud what that move costs.

Pressure, move, price: the map of this course

Every system in this roadmap is a mix of a few pressures. Each pressure has a standard move, and each move has a price. You will meet every row again in depth.

Pressure in the requirements Move Price you pay Deep dive
The same answer is requested again and again Cache it, in memory or at a CDN edge Stale answers; invalidation work Databases & Storage
Users are far away, or responses are big static files CDN Purges; per-GB bills Networking Basics
More requests than one server can compute Stateless servers behind a load balancer State must live somewhere shared Core Building Blocks
Reads outgrow one database Read replicas Replication lag, so some reads are stale Databases & Storage
Writes or data outgrow one database Partition (shard) the data Cross-shard queries; hot keys Databases & Storage
Slow or bursty work inside a request Queue + workers Work finishes later; jobs can run twice Messaging & Queues
Big files: photos, video, backups Object storage No queries or transactions across files Core Building Blocks
Two users must never get the same seat One authority decides (a transaction, a unique constraint) Slower, and less available during failures Key Concepts & Terminology
One post must reach millions of followers Fan-out on write, on read, or both Storage vs read latency Twitter/Instagram News Feed
One client sends far too many requests Rate limiting Some legitimate requests get rejected URL Shortener & Rate Limiter
Any part can fail Redundancy, timeouts, safe retries Money and complexity This lesson, then Advanced Topics & Final Prep

The roadmap itself runs like this: foundations, then the building blocks (databases, caches, queues), then the interview method (requirements, then design steps), then complete designs (a URL shortener and rate limiter, a news feed, video streaming, chat), and finally advanced topics and mock interviews.

This lesson teaches the three things everything else builds on: what the interview grades, the life of a request (the boxes in that table, drawn as a path), and whiteboard arithmetic. Then it uses them to grow a design and to plan for failure. It assumes you have built a web app and know what an HTTP request and a SQL query are.

Your turn: the page that cannot be shared

Same spike, same numbers, but the popular page is the logged-in "My cookbook" page: every reader sees their own saved recipes. What does the CDN still do for you, and what do you do about the rest?

Check your answer

The CDN still serves the images, CSS and JS, which are 95% of the bytes. The HTML alone is 5,000 Γ— 100 KB Γ— 8 = 4 Gbps, which fits the 10 Gbps link.

But this HTML cannot be shared: each reader's page is different, and each is requested once. Caching only pays when the same answer is requested again, and here it is not. So the app and the database must really do the work:

  • App: 5,000 renders/s Γ· 1,500 per server β‰ˆ 3.3 servers at full load, about 6 at a 60% target. That means stateless servers behind a load balancer.
  • Database: 20,000 queries/s against 8,000. The cheapest move is often not hardware. If one query with a join can replace the four, that is 5,000 heavier queries/s. Measure it, because the 8,000/s figure is for simple queries. Moving the database to its own machine, instead of sharing CPU with the app, adds headroom too. If that is still not enough, add read replicas.

The price: more servers to pay for and run. Replicas also bring replication lag, so a reader who just saved a recipe might not see it on the next page load. Name that, and say how you would handle it.

What the interviewer is actually grading

"Design a food delivery app" in 45 minutes has no answer key. Two strong engineers can draw different diagrams and both pass. What the interviewer grades is judgment: whether your choices follow from the requirements, and whether you know what each one costs.

Not a coding interview

Coding interview System design interview
The problem Precisely stated Deliberately vague; you narrow it
Success Correct output, good complexity Sound choices with stated trade-offs
Your time Mostly writing code Mostly talking and drawing
A "right answer" Usually exists Several good ones exist
Who drives The problem statement You, with the interviewer steering

Five signals, and what they sound like

Signal What strong candidates do What it sounds like
Scoping Turn the vague prompt into a few features and a few numbers "Let's cover posting, following and the home feed, and leave out messaging and ads. OK?"
Quantifying Estimate before choosing "That's about 17,000 feed reads per second at peak, too many to send to one database."
Trade-offs Name the price of every move "A 60-second cache means an edit takes up to a minute to show. For recipes, that's fine."
Failure Ask what happens when each box dies, slows down or gets a retry "If the cache node dies, the database sees ten times its load, so…"
Communication Drive, check in, adapt when steered "I can go deeper on the feed or on photo storage. Which is more useful to you?"

One sentence shape is worth rehearsing until it comes out automatically:

"I'll use X because [requirement]. It costs Y. If [requirement] changes to Z, I'd switch to W."

Interviewers push on exactly the last part. Expect follow-ups such as "what if traffic is 10Γ— higher?", "what happens when that node dies?" and "how would you know it's broken?" If you have already named the price, those questions are easy.

Most trade-offs pull between three goals that Martin Kleppmann names in Designing Data-Intensive Applications: scalability (coping with growth), reliability (working correctly when parts fail) and maintainability (people can run and change it). A design that chases the first two by adding ten services usually loses the third.

Requirements first, because they change everything

"Design a messaging app" could mean a consumer messenger for hundreds of millions of people, whose phones go offline for days and whose messages are end-to-end encrypted. It could also mean a notification feed for the 1,000 employees of one company. The first needs a whole lesson (Chat Application (WhatsApp)); the second is one server and a table. You cannot tell which until you ask.

Requirements come in two kinds:

  • Functional: what the system does. "Users send messages to a group; messages sent while you are offline arrive when you reconnect."
  • Non-functional: how well it does it: scale, latency, availability, durability, consistency, cost. "99% of messages arrive within a second; an accepted message is never lost."

The non-functional ones decide the architecture, so those are the ones to turn into numbers. This is the Scope phase of the interview (below). Requirements Gathering shows how to ask for them quickly.

Predict first: the prompt is "Design a leaderboard for a mobile game." Two candidates open like this.

  • A: "I'll keep scores in Redis sorted sets, stream score events through Kafka and store history in Cassandra."
  • B: "A few questions first. One global ranking, or rankings among friends? How many players, and how many score updates per second at peak? Must a player see their new rank the moment a match ends, or is a minute later fine?"

Which opening is stronger, and why can A's design be good and still score badly?

Check your answer

B. Each of B's questions changes the design:

  • Friends-only rankings involve a few hundred players each, and you can sort them when someone asks. A global ranking of millions of players needs a structure that stays sorted as scores arrive.
  • 500 updates per second fits one in-memory sorted set on one node. 500,000 per second is past what one node should carry, so you partition the scores and merge the top entries.
  • "Instantly" means updating the ranking inside the request. "Within a minute" lets you batch score events behind a queue.

A may end up with a perfectly reasonable design. But nothing A said ties a tool to a requirement, so the interviewer cannot tell luck from judgment, and A has no numbers to defend the choice when pushed.

The shape of the session

This course uses one interview framework throughout, in four phases: Scope β†’ Sketch β†’ Deep dive β†’ Wrap-up.

Phase What you produce 45-minute slot 60-minute slot
Intro and prompt Nothing yet 0–3 0–5
1. Scope Features, qualities, two or three numbers, assumptions 3–10 5–13
2. Sketch API, data model, and a diagram with one write and one read traced end to end 10–20 13–25
3. Deep dive The riskiest part designed properly, with its failure modes (one topic in 45 minutes, two in 60) 20–35 25–48
4. Wrap-up Summary, top risks, what changes at 10Γ— 35–40 48–55
Your questions Your questions for the interviewer 40–45 55–60

Everything in this lesson has a slot in it. Requirements and estimates belong to Scope, and only the numbers that decide something. The life of a request is what you draw and trace in the Sketch. Failure thinking fills the Deep dive and comes back in the Wrap-up. A longer slot buys depth, a second deep-dive topic, not a longer requirements phase. Interview Framework & Strategy teaches the framework, its checkpoints and what to do when you fall behind; Design Process Steps covers the Sketch and the Deep dive.

The life of a request

Every design in this course is a variation on one path. Learn it once and you have a map: when a prompt mentions a pressure, you will know which box it lands on.

 client (browser or app)
   β”‚  1. DNS: shop.example.com β†’ an address of the CDN
   β–Ό
 CDN edge ─────────── 2. static files and cacheable pages are answered here
   β”‚  3. everything else is forwarded to your origin
   β–Ό
 load balancer ────── 4. picks a healthy app server
   β”‚
   β–Ό
 app server (one of many identical, stateless copies)
   β”œβ”€β”€β–Ί cache ─────────────── 5. look here first
   β”œβ”€β”€β–Ί database ──────────── 6. on a miss: read here, then fill the cache
   β”‚      └── replicas          (copies of the database that serve reads)
   β”œβ”€β”€β–Ί object storage ────── big files: photos, PDFs, backups
   β”œβ”€β”€β–Ί queue ──► workers ─── 7. slow work, finished after the response
   β”‚
   β–Ό
 8. the response travels back along the same path

Without a CDN, step 1 returns the load balancer's address and steps 2 and 3 disappear.

Hop by hop

Hop Its job You add it when… Its price
DNS Turns a name into an IP address. Resolvers cache the answer for the record's TTL Always there; you choose the records and TTLs A change usually takes up to the TTL to reach everyone, and some clients hold on longer, so DNS is a slow failover lever
CDN Caches responses in many cities, and terminates connections close to users Users are far away, or bytes are big and shared Purges, staleness, per-GB bills. Personal data must never be cached as shared
Load balancer Spreads requests over servers; health checks take dead ones out You run more than one app server, for capacity or survival One more hop. It is a single point of failure itself unless it is redundant
App servers Run your code and keep no user state in memory Always; add copies as CPU runs out Sessions, carts and uploads must move to shared stores
Cache Keeps hot data in memory (Redis, Memcached) The same reads repeat and a little staleness is acceptable Stale reads until the TTL or an invalidation. A cold or dead cache sends the full load to the database
Database The durable source of truth Always The hardest box to scale: replicas for reads (lag), partitioning for writes
Object storage Stores whole files by key, cheaply and durably Files, media, backups No queries; the metadata about each file lives in a database
Queue + workers Holds jobs until workers take them Slow, bursty or retry-prone work the user need not wait for Results arrive later. Most queues deliver at least once, so a job can run twice

Deep dives: the load balancer, app tier and object storage in Core Building Blocks; caches and databases in Databases & Storage; queues in Messaging & Queues. Three details on the path trip people up.

Stateless means "any server can take any request". If server 2 keeps your cart in its own memory, the load balancer must keep sending you to server 2 (a sticky session), and a restart or crash of server 2 empties your cart. Keep that state in a shared store and the servers become interchangeable: you can add them, kill them and deploy them one at a time. Sticky sessions are not forbidden, but they have a price: uneven load, and lost state when that server goes.

Cache-aside: the app talks to both. On a read, the app asks the cache. On a miss, the app reads the database and writes the answer into the cache with a TTL (time to live). The cache never talks to the database in this pattern, which is why the diagram draws two separate arrows from the app server. Other strategies and invalidation are in Databases & Storage.

202 Accepted, not 200 OK. When work goes on a queue, the honest answer is "accepted, not done yet", and HTTP has a status code for exactly that: 202. The client then polls, gets notified, or never needs to know.

Two transport facts you will use constantly

  • HTTP runs over TCP, and HTTPS adds TLS in between. HTTP/1.1 and HTTP/2 use TCP. HTTP/3 runs over QUIC, which runs over UDP and combines the transport and encryption handshakes. A new TCP connection costs one round trip before the request can leave, and a TLS 1.3 handshake costs one more (RFC 8446). QUIC does both in one round trip (RFC 9000).
  • Plain UDP is for real time. Voice and video calls and multiplayer games use UDP, because a late packet is worthless and waiting for a retransmission causes a stall. On-demand video, such as a film or a recorded clip, is delivered over HTTP, so over TCP or QUIC.

Networking Basics covers DNS, HTTP, CDNs, TCP, UDP and QUIC properly.

Trace: one request from Frankfurt

Your shop's servers are in Virginia, and a customer in Frankfurt opens a product page for the first time. The assumptions:

  • A Frankfurt–Virginia round trip takes about 90 ms. The straight-line distance, about 6,500 km, sets a floor near 65 ms in fiber, and real routes are longer.
  • The DNS answer is not cached yet, and the lookup takes 30 ms.
  • The server's own work (load balancer, app code, one cache hit) takes 8 ms.

Predict first: of the roughly 300 ms this request takes, how much is your server's work? Which single change would save the most?

Check your answer
Step Straight to Virginia Through a CDN edge in Frankfurt
DNS lookup 30 ms 30 ms
TCP handshake (1 round trip) 90 ms 10 ms (the edge is close)
TLS 1.3 handshake (1 round trip) 90 ms 10 ms
Request and response (1 round trip) 90 ms 10 ms, plus 90 ms from the edge to Virginia over a connection it keeps open
Server work 8 ms 8 ms
Total 308 ms 158 ms

Your server's work is 8 ms, under 3%. Three trips across the Atlantic take 270 ms. Making all of the server's work twice as fast would save 4 ms. The big savings come from fewer long round trips, and for this first visit the CDN edge saves the most:

  • A nearby CDN edge does the handshakes 10 ms away and reuses a warm connection to Virginia. That gives 158 ms, even though the page itself is dynamic. A static file cached at that edge takes about 61 ms.
  • Reusing the same connection skips DNS and both handshakes: 98 ms. That helps only the requests that follow, not this first one.
  • HTTP/3 merges the two handshakes into one round trip: 218 ms, even straight to Virginia.

Why it works. On an idle system, a request's latency is roughly round trips Γ— distance, plus work. Physics sets the first term. Light in fiber covers about 200,000 km per second, so every 1,000 km between client and server adds about 10 ms to each round trip. You cannot speed up light. You can only move the answer closer, or cross the distance fewer times. Under heavy load a third term appears, waiting in queues, and it grows sharply as servers approach full utilization. Key Concepts & Terminology covers that and tail latency.

Whiteboard arithmetic

Estimates exist to choose boxes. "32 petabytes of photos" tells you the photos cannot live in the database. "125 Gbps at peak" tells you a CDN is not optional. An estimate that does not change a decision is time you could have spent designing, so after each number, say what it decides.

Numbers worth memorizing

These are the classic "latency numbers every programmer should know" from Peter Norvig and Jeff Dean, rounded (the list). Modern hardware beats some of them (NVMe SSDs and newer datacenter networks are faster), but you reason with the ladder of orders of magnitude, and that has not changed.

Operation Time
Main memory reference 100 ns
Read 4 KB at random from an SSD 150 Β΅s
Read 1 MB sequentially from memory 250 Β΅s
Round trip inside one datacenter 0.5 ms
Read 1 MB sequentially from an SSD 1 ms
Seek on a spinning hard disk 10 ms
Read 1 MB sequentially from a spinning disk 20 ms
Packet from California to the Netherlands and back 150 ms

Three consequences you will use:

  • Memory is about 100,000Γ— faster than a disk seek (100 ns against 10 ms). A database whose hot data fits in RAM behaves very differently from one that has to seek.
  • A cache hit on another machine costs a network round trip, about 0.5 ms, not 100 ns. Ten sequential cache calls in one request add about 5 ms, so batch them into one round trip.
  • Crossing an ocean costs more than almost anything you do inside the datacenter. Count the long round trips first.

Units that bite

Fact Value Whiteboard note
Seconds in a day 86,400 Rounding to 10⁡ makes a per-second rate about 14% low. Fine, but say so
1 million per day about 12 per second
1 billion per day about 11,600 per second
Seconds in a year about 31.5 million about 3 Γ— 10⁷
1 byte 8 bits Network speeds are quoted in bits, file sizes in bytes
1 Gbps 125 MB/s
K, M, G, T, P steps of 10³ 2¹⁰ = 1,024 is 2.4% more per step, so about 7% at G and 13% at P; still irrelevant at this precision

Formulas you will reuse

You want Formula Watch out for
Average rate events per day Γ· 86,400 Per day vs per second
Peak rate average Γ— peak factor Assume 2–3Γ— for ordinary daily cycles if nobody says; a launch or a TV mention can be 100Γ—
Storage new items per day Γ— bytes each Γ— days kept Γ— copies Retention and copies. Managed object storage bills the bytes you store and handles its own redundancy
Bandwidth requests per second Γ— bytes each Γ— 8 Bits vs bytes
Servers peak rate Γ· (one server's rate Γ— target utilization), plus spares Planning for servers that run at 100%
Cache memory hot items Γ— bytes each Caching everything instead of the hot set

Worked example: a photo app

The interviewer offers these numbers. They are assumptions for the exercise, not any real company's figures.

Assumption Value
Daily active users 50 million
Feed opens per user per day 10, each showing 6 photos at about 150 KB
Uploads 10% of users post one photo a day: 3 MB original plus 0.5 MB of resized copies
Metadata per photo 1 KB (owner, caption, time, file keys)
Peak factor 3Γ— the daily average
Plan for 5 years

On the whiteboard:

photo views   50M Γ— 10 Γ— 6  = 3B/day       β‰ˆ 35k/s avg    β‰ˆ 104k/s peak
feed calls    50M Γ— 10      = 500M/day     β‰ˆ 5.8k/s avg   β‰ˆ 17k/s peak
uploads       5M/day                       β‰ˆ 58/s avg     β‰ˆ 174/s peak
photo bytes   5M Γ— 3.5 MB   = 17.5 TB/day  β‰ˆ 6.4 PB/year  β‰ˆ 32 PB in 5 years
egress        35k/s Γ— 150 KB Γ— 8 bits      β‰ˆ 42 Gbps avg  β‰ˆ 125 Gbps peak
ingress       58/s Γ— 3 MB Γ— 8 bits         β‰ˆ 1.4 Gbps avg (egress is β‰ˆ 30Γ— this)
metadata      5M Γ— 1 KB     = 5 GB/day     β‰ˆ 1.8 TB/year

Now the part that earns the points: what each number decides.

Number Decision
32 PB of photos Photos go to object storage. The database holds only each photo's 1 KB of metadata and its file key
125 Gbps at peak, 30Γ— the upload traffic Images are served through a CDN. The origin must not be what every view hits
17k feed calls/s at peak Stateless app servers behind a load balancer. If one handles 1,000/s (measure it), about 29 at a 60% target
600 views per upload Read-heavy, so cache feed and metadata reads. The last two days of posts, 10M Γ— 1 KB β‰ˆ 10 GB, fit on one cache node
174 metadata writes/s at peak (58 average), 1.8 TB/year One database primary with a standby handles this for years. Plan partitioning by user before the data outgrows one server

In the room you say only the numbers that decide something. Say it aloud in about this many words: "Roughly 17,000 feed requests a second at peak against under 200 uploads a second, so it's heavily read-dominated. Photos add 17 terabytes a day, over 30 petabytes in five years, so they go to object storage behind a CDN, which carries about 125 gigabits at peak. Metadata is small, under two terabytes a year, so I'd start with one database, a standby and a cache in front." Then move on. Precision beyond one or two significant figures rarely changes a decision.

The classic slips

Slip In this example Off by
Reading per day as per second "3 billion views a second" 86,400Γ—
Mixing bytes and bits "42 GB/s of egress" (it is 42 Gbps, about 5.2 GB/s) 8Γ—
No peak factor Provisioning for 35k views/s 3Γ— short at peak
Forgetting copies A self-run cluster keeping 3 replicas needs 3Γ— the raw disk 3Γ—
Forgetting retention Quoting "17.5 TB" as the total 1,825Γ— (five years of days)
Over-precision "34,722.2 views per second" Wastes a minute; round

Your turn: swap photos for video

Product wants 15-second video clips instead of photos, encoded at 2 Mbps. The numbers of uploads and views stay the same, and every clip in the feed autoplays in full. Which number changes most, and what do you do about it?

Check your answer

A clip is 15 s Γ— 2 Mbps = 30 megabits = 3.75 MB. That is about the size of a photo with its resized copies (3.5 MB), so storage barely moves: about 19 TB a day, 34 PB over five years.

Egress explodes. Each view now moves 3.75 MB instead of 150 KB, 25 times more: about 1 Tbps on average and 3 Tbps at peak. The CDN becomes the biggest line on the bill. The cheapest fixes are product and encoding decisions, not infrastructure: do not autoplay every clip, or autoplay feed clips at a capped low bitrate and switch to full quality only when the viewer taps in. Separately, the player needs several renditions of each clip so it can pick one that matches the viewer's bandwidth (adaptive bitrate). Storing those renditions raises storage, perhaps 2Γ—, and transcoding each upload into them is slow work, so it goes on a queue. YouTube / Netflix Streaming does this properly.

Earn every box

A single server is a real design

A school's homework portal serves 800 students. Say each makes 30 requests on a busy day: 24,000 requests, or 0.3 per second on average. Even if every student submits in the last hour before a deadline, that is about 7 per second. One server with a managed database and backups handles it with plenty of room to spare.

In an interview, drawing Kubernetes, Kafka and six microservices for this is not ambition. It is a design that costs more, fails in more ways and needs more people, to solve problems that do not exist. Every box should be the answer to a number.

"Monolith" means one deployable unit, not one machine: a monolith can run as 20 identical copies behind a load balancer. You split out a separate service when one part must scale, deploy or fail independently, or when teams keep blocking each other. You pay for it with network calls that can fail, harder debugging and more to operate. Shopify has described deliberately keeping a modular monolith, with enforced boundaries inside one codebase, instead of moving to microservices (Shopify Engineering, 2019). Advanced Topics & Final Prep covers microservices.

Traffic grows 10Γ—: who breaks first?

Your shop API runs on four stateless app servers, one cache node, and one Postgres primary with a standby. Here is peak load today, from your dashboards:

Tier Today at peak
App servers (4) 30% CPU
Cache node 5% CPU; the hot set uses 12 of 64 GB
Database primary 35% CPU: 25 points from cache-miss reads (hit rate 90%), 10 points from writes

A new market launches, and peak traffic will be 10Γ— higher with the same mix of requests.

Predict first: which tier breaks first, and which is hardest to fix?

Check your answer

Multiply by ten:

Tier At 10Γ— Verdict
App servers 300% of today's fleet Broken but easy: add servers. That is 12 servers' worth at full load, about 20 at a 60% target
Cache node 50% CPU; memory unchanged if the hot set is Fine for now. Add a replica so one crash does not empty it
Database 250 points of reads + 100 of writes = 350% Broken, and hard

The database is worst, and it is the one you cannot fix by adding identical copies. Split its load by cause:

  • Reads (250 points). Raise the cache hit rate from 90% to 99%, with a bigger cache or longer TTLs if the hot set allows it. Misses fall from 10% of reads to 1%, ten times fewer: 25 points. Read replicas are the other move, and their price is replication lag.
  • Writes (100 points). That alone fills a whole primary, and replicas do not help: every replica must apply every write too. The cheap move is a primary twice as big, which puts writes at 50%, plus removing needless writes (for example, updating "last seen" on every request). The next doubling needs partitioning: splitting the data by customer across several primaries.

After the hit-rate fix and a 2Γ— bigger primary, the database runs at (25 + 100) Γ· 2 β‰ˆ 63% at peak. That buys time, not forever, and a strong candidate says so.

Why the stateful tier is the hard part

App servers are easy to multiply because they hold nothing: any copy can serve any request, so ten copies do ten times the work. A database is its state. Copies must agree on what is true, so either every write reaches every copy (replication), or each node owns a different slice of the data (partitioning). Both have a price: lag, coordination, and queries that span slices. That is why the database is usually the first hard bottleneck in a growing web service, and why so much of this course is about data.

🧠 Name the bottleneck before the move. "At 10Γ— the primary is at 350%: 250 points from reads, 100 from writes. The cache fixes the reads. Writes need a bigger primary now and partitioning by customer later, because replicas don't take writes." That is the difference between adding boxes by instinct and adding them because a number demanded it.

Scale up, then out

Buying a bigger machine (vertical scaling) needs no code changes, and it is a perfectly good first move. Its limits are real, but they are mostly not about price. In the cloud, the hourly price within one instance family grows roughly in proportion to size, so one 64-vCPU machine costs about the same as eight 8-vCPU machines of that family. The limits are these:

  • There is a biggest machine. Past it, you must split.
  • One machine is one failure. A bigger box is no more redundant than a small one, so you need a standby anyway.
  • Resizing usually means a restart or a failover.

Adding machines (horizontal scaling) has no hard ceiling for stateless servers, and it gives you redundancy along the way. For the database it means replicas and partitioning, with the prices above. So keep the app tier stateless from day one, so you can scale it out the moment a number demands it. Scale the database up while a bigger machine exists and the step is easy, and plan partitioning before you run out of bigger machines. Key Concepts & Terminology treats scalability properly.

Failure is the normal case

Netflix has described how its move to the cloud began: in August 2008 a major database corruption stopped it from shipping DVDs for three days (Netflix, 2016). It later built Chaos Monkey, which randomly terminates production instances so that engineers must build services that survive losing one. Its schedule runs on weekdays, between 9:00 and 15:00 by default, which keeps the breakage to hours when engineers are at work to respond (configuration docs).

At scale, something is always broken: a disk, a server, a network link, a dependency. The design question is never whether a box fails but what the user sees when it does.

Three questions for every box

Point at each box on your diagram and ask:

  1. What if it dies? Does everything stop, or does the system degrade?
  2. What if it is slow? Do callers wait forever, or time out and fall back?
  3. What if a request to it is retried? Does the retry repeat a side effect, such as a second charge or a second email?

Failure modes along the path

Event What happens Move What you give up
An app server dies Its in-flight requests fail; health checks take it out after a few failed probes (seconds to a minute, depending on settings) Several stateless servers and spare capacity (N+1) Paying for idle headroom
The load balancer dies Everything behind it is unreachable A redundant pair, or a managed load balancer spread across zones Money; one more thing to configure
The cache dies or restarts empty Every read misses. The database gets many times its normal load and can collapse Cache replicas, one database fetch per key at a time, throttled misses, warming before serving Memory and complexity
The database primary dies Writes fail until a standby takes over. With asynchronous replication, the last committed writes can be lost A standby with automatic failover; synchronous replication if no loss is acceptable Seconds to minutes of failed writes; with synchronous replication, every commit waits for the standby
A dependency gets slow Threads pile up waiting, and the whole service stalls Timeouts, a circuit breaker, a fallback That feature degrades (for example, generic recommendations)
Clients retry after timeouts Load multiplies just when you are overloaded; side effects repeat Exponential backoff with jitter, retry limits, idempotency keys Some users wait longer
A whole zone or region fails Everything in it is gone Spread across zones by default; several regions when required Cost, cross-region latency, harder consistency

The PostgreSQL documentation states the asynchronous-replication risk plainly: committed transactions that had not reached the standby are lost on failover (docs). Retries deserve their own warning: if three layers each try a call four times (once plus three retries), one user request can become 4 Γ— 4 Γ— 4 = 64 calls at the bottom layer. Retry at one layer, back off, and add random jitter so clients do not retry in lockstep.

The retry that charges twice

The third question catches the most expensive bugs. Your checkout calls a payment provider. The provider charges the card, but the reply is lost: the network drops, or your HTTP client gives up after 5 seconds while the provider finishes after 6. A timeout means "I don't know", not "it failed". If you retry, you may charge twice. If you don't, you may not have charged at all.

The fix is an idempotency key: a unique ID for the operation (this checkout), created once and sent with every retry. A provider that supports keys charges at most once per key, for as long as it remembers the key. Stripe's API, for example, accepts an Idempotency-Key header and returns the saved result when the same key comes back (docs). A toy provider shows the difference:

class PaymentProvider:
    """A toy provider: the charge really happens, but a reply can get lost."""
    def __init__(self, lost_replies=0):
        self.charges = []             # what the bank actually did
        self.replies = {}             # idempotency key -> original reply
        self.lost_replies = lost_replies

    def charge(self, amount, key=None):
        if key is not None and key in self.replies:
            return self.replies[key]  # a retry: replay the first reply
        reply = f"charge-{len(self.charges) + 1}"
        self.charges.append(amount)
        if key is not None:
            self.replies[key] = reply
        if self.lost_replies > 0:     # charged, but the caller never hears back
            self.lost_replies -= 1
            raise TimeoutError("no reply")
        return reply


def pay(provider, amount, key=None, attempts=3):
    for _ in range(attempts):
        try:
            return provider.charge(amount, key=key)
        except TimeoutError:
            pass                      # outcome unknown, so try again
    raise RuntimeError("gave up")


plain = PaymentProvider(lost_replies=1)
pay(plain, 40)
print("no key:  ", plain.charges)     # [40, 40]: charged twice

keyed = PaymentProvider(lost_replies=1)
pay(keyed, 40, key="order-7731")
print("with key:", keyed.charges)     # [40]: charged once

Your own service needs the same guarantee, because your clients retry too, and people double-click "Pay":

  1. Reserve the key first. Insert a row with the key and the status pending into a table where the key is UNIQUE. A duplicate request collides with that row. If the row says done, return the stored result. If it says pending, answer "still processing", or carry on at step 2, which is safe because of the key.
  2. Call the provider with the same key.
  3. Record the result on that row, then return it.
  4. If you crash between steps 2 and 3, the row stays pending. A retry, or a reconciliation job, calls the provider again with the same key, and the provider replays the original result instead of charging again. That works only within the provider's key-retention window (Stripe may prune keys once they are 24 hours old). After that, look the charge up first, for example by the order ID you stored in its metadata, before trying again.

The price is a table, a pending state and a reconciliation job. The classic broken version looks for an existing row first and inserts one only after charging, with no unique constraint. Two concurrent requests both see "no row" and both charge; a crash after the charge leaves no row, so the retry charges again. Queues redeliver jobs too, and the same idea makes a queue consumer safe to run twice; Messaging & Queues covers that.

Not all data needs the same guarantee

During a failure, some reads can only be answered from a copy that might be out of date. Whether that is acceptable depends on the operation, not the whole system. A like count a few seconds behind hurts nobody, so serve it from any replica. "Is seat 14C still free?" must be decided by one authority, a single primary with a transaction or a unique constraint, even if that means refusing to sell seats while the authority is unreachable. Decide per operation, and say which operations get which treatment. Key Concepts & Terminology states the CAP theorem and the consistency models precisely.

Your turn: the region goes dark

Your shop runs in one cloud region, spread across three availability zones. The interviewer asks: "What happens if the whole region goes down? Would you go multi-region?" Answer in two or three sentences, as you would in the room.

Check your answer

Start from the requirement, not the technology. Two numbers frame it: how long you may be down (the recovery time objective, RTO) and how much recent data you may lose (the recovery point objective, RPO).

A strong answer: "With three zones we survive losing a datacenter. A whole-region outage is rare, and for a shop an RTO of about an hour and an RPO of a few minutes is usually acceptable, so I'd keep backups and replicated data in a second region, plus a tested runbook to rebuild there, rather than run two live regions. If the business needs an RTO of minutes and almost no data loss, I'd run a warm second region. That can approach double the cost, and it means either accepting a small loss window from asynchronous cross-region replication, or adding a cross-region round trip to every write with synchronous replication."

The follow-up to expect is "how do you know the runbook works?" The answer: you rehearse it regularly.

Final round: no label on the problem

Real prompts don't say which move they want. Start from the requirements, put numbers on them, find the pressure, then pick the move and its price.

Challenge 1: results day

A national exam board publishes results for 3 million students at 09:00. In past years 60% checked in the first ten minutes, and the first minute was three times busier than the ten-minute average. Each student logs in and sees only their own result. Results never change after publication.

Before peeking, answer these:

  1. What is the peak request rate, and how many app servers does it take if one handles 1,000 requests per second?
  2. Where do the results live at 09:00, and why?
  3. Which caching mistake here would be a scandal rather than a slowdown?
  4. Why is autoscaling alone not enough?
Hint

The load is a wall, not a curve. Every answer is personal, but none of them change.

Check your answer
  1. 1.8 million checks in 600 s is 3,000 per second on average, and 9,000 per second in the first minute. At 1,000 per server that is 9 servers at full load, or 15 at a 60% target, plus whatever login needs.
  2. Results are fixed, so compute them before 09:00 and load them into a key-value store or cache keyed by student ID. At about 2 KB per result (an assumption), 3 million Γ— 2 KB β‰ˆ 6 GB fits in memory on one node, plus a replica. At 09:00 every lookup is a memory read, not a database query.
  3. Caching a personal results page at the CDN under a shared key. The page shell, CSS and JS can be cached. The result itself must be fetched per student and marked as not cacheable by shared caches. Get this wrong and students see each other's results.
  4. New servers usually take minutes to start and warm up, and the spike is over in minutes. Scale out before 09:00 by pre-provisioning the 15 or so servers, and keep a waiting-room page ("you're in the queue") ready in case the estimate is wrong.

Challenge 2: the overnight report run

Every night at 01:00, a B2B product generates a PDF report for each of its 20,000 customers and emails it. One report takes 3 seconds of CPU. Reports must arrive by 08:00, and no customer should get the same report twice. How many workers do you need, and what happens when one crashes halfway?

Check your answer

Size it. 20,000 Γ— 3 s = 60,000 CPU-seconds, about 16.7 CPU-hours. In a 7-hour window that needs at least 2.4 busy cores. With 8 workers the run takes about 2 hours, which leaves time to redo failures.

Shape it. A scheduler puts one message per customer on a queue. A worker takes a message, renders the PDF, stores it in object storage, sends the email and only then acknowledges the message. If a worker crashes, its message is never acknowledged, so the queue hands it to another worker. Nothing is silently skipped.

The duplicate. Because of redelivery, a worker can crash after sending the email but before acknowledging, and the report goes out again. So record progress per customer and date: mark "sending" before the send and "sent" after it, and skip customers already marked sent. That shrinks the duplicate window to a crash between those two marks. As stated, the requirement goes further: closing that window needs an email provider that deduplicates on a key you supply, such as customer plus date. If no such provider is available, tell the interviewer that at-least-once delivery means a rare duplicate, and ask whether that is acceptable.

System design foundations in five minutes: find what breaks first, pick the cheapest move, say what it costs

Cheat sheet: pressure β†’ move β†’ price

When the requirements say… Reach for Say the price
"the same thing is read constantly" A cache or a CDN "Readers can see data up to one TTL old."
"users worldwide", "big files" A CDN; object storage for the files "Edits need a purge, and we pay per GB served."
"more traffic than one server" Stateless servers + a load balancer "Sessions and carts move to a shared store."
"reads outgrow the database" Cache first, then read replicas "Replicas lag; a user's own reads may need the primary."
"writes or data outgrow the database" A bigger primary, then partitioning "Cross-partition queries get harder, and hot keys need care."
"slow work", "bursts" Queue + workers "It finishes later, and a job can run twice."
"must never double-charge or double-send" Idempotency keys, reserved before the side effect "A pending state and a reconciliation job."
"must stay up if X dies" Redundancy, timeouts, failover "Idle capacity, and seconds to minutes of failover."
nothing that needs it Nothing. One box "Simpler to run. I'll add parts when a number demands it."

Before moving on, pick a product you used today (a food-delivery app, your bank, a video site) and explain it aloud as in an interview, in five minutes. Scope it with three clarifying questions and two or three numbers, each with the decision it drives. Sketch the path from DNS to database, tracing one read and one write. Go deep on one failure and the price of its fix. Then wrap up in two sentences: the top risk, and what changes at 10Γ—. If you can do that without notes, the rest of this course adds depth to a map you already hold.

Next: Key Concepts & Terminology states scalability, latency, availability and consistency precisely, and Networking Basics covers DNS, HTTP, CDNs, TCP, UDP and QUIC. Then Core Building Blocks opens up each box on the path.