Why start with a story
System design interviews ask you to draw an architecture for something like a URL shortener, a chat app or a news feed. Beginners often freeze because the final diagrams in books look like a city map: load balancers, caches, queues, replicas, shards, CDNs. It is hard to remember thirty boxes when you do not know why any of them exist.
The good news is that every one of those boxes was added to solve one specific problem. Real systems do not start as a city map. They start as one computer, and they grow one problem at a time. If you understand the problem each component solves, you never have to memorise a diagram again. You can rebuild it from first principles in an interview, and you can explain why each box is there, which is exactly what interviewers probe.
This lesson follows one made-up app, Recipebox, from its first hundred users to tens of millions. At every step you will see:
- The symptom: what users or engineers notice going wrong.
- The fix: the component or technique that solves it.
- The new trade-off: what the fix costs you, because nothing is free.
- The diagram: a text drawing that grows each time.
You will meet the same Recipebox in the next three lessons of this beginner path, so the numbers and names will stay familiar.
Meet Recipebox
Recipebox lets people post recipes with a photo, browse a feed of new recipes, search by ingredient, leave comments and get an email when someone comments on their recipe. It is simple enough to understand in one sentence, but it needs reads, writes, images, search, and background work, which is everything a typical interview problem needs.
Throughout, we use small rounded numbers. A useful rule of thumb that we will use repeatedly: a typical Recipebox user makes about 30 requests per day (page loads, API calls, image fetches from our servers). There are about 86,400 seconds in a day, and traffic at the busiest hour is roughly 3 times the daily average. These are assumptions, not laws, and the estimation practice lesson shows how to choose them.
Stage 0: Everything on one server
On launch day Recipebox has a few hundred users, mostly friends of the founders. The entire system runs on one rented virtual machine.
+-----------+
users -->| Server |
| web app |
| database |
| photos |
+-----------+
Here is what happens when a user opens the home page:
- The browser asks DNS (the Domain Name System, the internet's phone book) for the address of
recipebox.example. DNS answers with the server's IP address, for example203.0.113.10. - The browser sends an HTTP request to that address.
- The web application code runs, asks the database on the same machine for the ten newest recipes, builds a page and sends it back.
- The browser then requests each recipe photo, which the same server reads from its local disk.
With 1,000 daily users making 30 requests each, that is 30,000 requests per day, or about 0.35 requests per second on average and roughly 1 per second at peak. A single modest machine handles that without noticing.
Why this is the right design at this size: one server is cheap, easy to deploy and easy to debug. Every log is in one place. There is no network between the app and the database, so queries are fast. Adding a load balancer or a cache now would be premature optimisation: you would pay for complexity before it buys you anything.
What is wrong with it: everything is a single point of failure (SPOF), meaning one component whose failure takes the whole system down. If the machine reboots, the disk fills or a bad deploy crashes the process, Recipebox is offline, and if the disk dies, every recipe and photo is gone unless you have backups.
Interview tip
Starting with one box in an interview is not naive if you say why: "At this traffic one server is fine; I will show what breaks first as we scale." Interviewers like candidates who justify complexity rather than drawing every component on minute one.
Stage 1: Move the database to its own machine
Symptom. Recipebox gets featured on a food blog and reaches 50,000 daily active users (DAU). That is about 1.5 million requests per day, roughly 17 per second on average and about 52 at peak. The site is still up, but pages slow down every evening. Engineers notice that the app process and the database are fighting over the same CPU and memory. The database wants lots of RAM to keep hot data in memory; the app wants CPU to render pages. Every app deploy also restarts the machine the database lives on.
Fix. Put the database on its own server. The app now talks to it over the private network.
+-----------+ +-----------+
users -->| App server|------->| Database |
| + photos | | server |
+-----------+ +-----------+
Why it helps:
- Each machine can be sized for its job. Databases like lots of memory and fast disks; app servers like CPU.
- You can restart or redeploy the app without touching the database.
- You can now scale the two tiers independently, which we will need soon.
New trade-off. Every query now crosses the network. Inside one data centre a round trip is roughly half a millisecond, which is tiny compared with a page render, but a page that makes 50 queries in a loop now spends 25 ms just on network hops. This is the first time you notice that how many calls you make matters, not only how fast each one is. You also now have two machines to monitor, patch and back up.
Stage 2: Buy a bigger machine (vertical scaling)
Symptom. At 1 million DAU, Recipebox serves about 30 million requests per day. That is roughly 350 requests per second on average and about 1,000 at the evening peak. The app server's CPU sits at 95% every evening and requests start timing out.
Fix. The easiest fix is vertical scaling (also called scaling up): replace the server with a bigger one, with more CPU cores, more RAM and faster disks. No code changes are needed. You do the same for the database server.
+=============+ +=============+
users -->| BIG app |------->| BIG |
| server | | database |
+=============+ +=============+
Vertical scaling is underrated. Modern machines with dozens of cores and hundreds of gigabytes of RAM can carry a surprising amount of traffic, and for databases in particular, scaling up is often the simplest first move because it keeps all the data in one place.
New trade-offs. Vertical scaling hits three walls:
- The ceiling. There is a biggest machine you can rent. When you outgrow it, there is no next size.
- The cost curve. Bigger machines often cost more than proportionally. Doubling cores can more than double the price at the high end.
- Still a single point of failure. A bigger box fails just as completely as a small one. Upgrading it usually means downtime too.
At this point it is worth separating two words that interviewers like to test. Performance is how fast the system serves one user. Scalability is whether it stays fast when you add more users and more resources. Recipebox has a scalability problem: each request is quick, but the system cannot absorb more of them. The building blocks lesson defines these and related pairs (latency vs throughput, availability vs consistency) in detail.
Stage 3: A load balancer and many stateless app servers
Symptom. The big app server is at its limit again, and a hardware fault took the site down for 40 minutes last week. Users noticed.
Fix. Instead of one big app server, run several ordinary ones and put a load balancer in front. A load balancer is a component that receives every incoming request and forwards it to one of several backend servers, skipping any that fail a health check (a small periodic request, such as GET /health, that confirms a server is alive). This is horizontal scaling (scaling out): adding more machines instead of bigger ones.
+-------------+
+--->| App server 1|---+
users --> +----+ | +-------------+ | +-----------+
| LB |-+--->| App server 2|---+--->| Database |
+----+ | +-------------+ | +-----------+
+--->| App server 3|---+
+-------------+
DNS now points at the load balancer, not at any one app server. If server 2 dies, the load balancer's health check fails for it within seconds and traffic flows to servers 1 and 3. To handle more traffic, add server 4.
The catch: app servers must be stateless
On a single server, Recipebox kept each user's login session in the server's memory. With three servers, that breaks. A user logs in on server 1, the next request goes to server 2, and server 2 has never heard of them, so they are logged out.
A stateless server keeps no data about a user between requests in its own memory or disk. Every request carries, or can look up, everything needed to serve it. To make Recipebox's app servers stateless:
- Move sessions to a shared store, such as an in-memory key-value database that all servers can reach, or
- Use signed tokens (for example, a signed cookie) that the client sends with every request so any server can verify who the user is.
- Move uploaded photos off the app server's disk (we will do this properly in Stage 7).
Once servers are stateless, they are interchangeable. You can add, remove or replace them at any time, and the load balancer can send any request anywhere.
Common mistake
"Sticky sessions" (the load balancer always sends a user to the same server) look like a shortcut, but they bring the old problem back: when that server dies, its users lose their sessions, and load spreads unevenly. Use them only as a stopgap, and say so in an interview.
New trade-offs:
- The load balancer itself is now a single point of failure. Production setups run a pair (one active, one standby, or both active) or use a managed load balancer that is already redundant.
- You now need to deploy code to many servers, which brings questions about rolling deploys and mixed versions.
- The database has quietly become the bottleneck: three app servers send three times the queries to one database.
For load-balancing algorithms (round robin, least connections, hashing) and L4 vs L7 balancers, see load balancing.
Stage 4: Database replication with read replicas
Symptom. At 10 million DAU, Recipebox handles about 300 million requests per day, roughly 3,500 per second on average and about 10,400 at peak. App servers scale out easily, but the database is at 90% CPU and slow queries pile up. Looking at the traffic, engineers find that about 90% of requests are reads (viewing recipes, browsing the feed) and only 10% are writes (posting, commenting, liking).
Fix. Use replication: keep copies of the database on several machines. One machine, the primary (older texts say master), accepts all writes. It streams every change to one or more read replicas (older texts say slaves), which apply the same changes and serve read queries.
writes
+-------------+ ---------> +-----------+
users --> +----+ | App servers | reads | Primary |
| LB |-->| (1..N) |---+ +-----------+
+----+ +-------------+ | | replicates
| v
| +-----------+ +-----------+
+-->| Replica 1 | | Replica 2 |
+-----------+ +-----------+
The app sends INSERT, UPDATE and DELETE to the primary and most SELECT queries to a replica. Since 90% of traffic is reads, adding replicas multiplies read capacity. Replicas also help availability: if the primary dies, one replica can be promoted to become the new primary, a process called failover.
New trade-offs:
- Replication lag. Replicas apply changes a little after the primary, usually milliseconds, sometimes seconds under load. A user posts a recipe, the page reloads from a replica that has not caught up yet, and the recipe seems to be missing. A common fix is read-your-own-writes: for a few seconds after a user writes, send that user's reads to the primary.
- Writes do not scale. Every write still goes to one primary. Replicas add read capacity, not write capacity.
- Failover is tricky. If replication is asynchronous (the primary confirms a write before replicas have it), a crash can lose the last few writes. If it is synchronous (the primary waits for a replica), writes get slower.
This is your first taste of the trade-off between consistency (every read sees the latest write) and availability and speed. A deeper treatment is in scalability and databases at scale.
Stage 5: Add a cache
Symptom. Even with replicas, the same queries run over and over. The ten most popular recipes are viewed millions of times a day, and each view runs the same SELECT with the same joins. Replicas are expensive and each one still reads from disk for data that does not fit in its memory.
Fix. Put a cache in front of the database. A cache is a fast store, usually kept in memory, that holds copies of data that is expensive to fetch or compute. Recipebox uses the most common pattern, cache-aside (also called lazy loading):
- The app receives a request for recipe 42.
- It looks in the cache for the key
recipe:42. - Cache hit: the value is there; return it. No database query.
- Cache miss: read recipe 42 from the database, store it in the cache with an expiry time (a TTL, time to live, such as 10 minutes), and return it.
+-------------+ +---------+
users --> +----+ | App servers |----->| Cache | (check first)
| LB |-->| (1..N) | +---------+
+----+ +-------------+
| on miss / writes
v
+-----------+ +-----------+
| Primary |---->| Replicas |
+-----------+ +-----------+
Let us estimate the effect. At peak, about 9,400 of the 10,400 requests per second are reads. If the cache answers 80% of them (an 80% hit ratio), only about 1,900 reads per second reach the database. That is roughly a fivefold drop, often enough to remove several replicas.
Reading from memory takes microseconds, while a database query that touches disk can take milliseconds. Caches therefore improve both latency (how long one request takes) and load on the database.
New trade-offs:
- Stale data. The cache may hold an old copy. If an author edits a recipe and the cache still holds the old version for 10 minutes, readers see the old one. You either accept that (with a short TTL) or delete the cache entry when the data changes (invalidation), which is easy to get subtly wrong.
- Cold starts. After a restart the cache is empty, and every request misses at once. The database can be overwhelmed. This is sometimes called a thundering herd.
- Memory is limited. You cache the hot subset, not everything, and you need an eviction policy, such as least recently used (LRU), to decide what to drop.
Cache patterns (write-through, write-behind, refresh-ahead) are compared in the building blocks lesson and in depth in caching.
Stage 6: A CDN for static assets
Symptom. Recipebox has users in India, Europe and North America, but all servers are in one data centre in Mumbai. Users in Toronto say the site feels sluggish. Looking at a page load, engineers see that the HTML arrives reasonably quickly, but the page also loads about 2 MB of JavaScript, CSS, fonts and icons, each fetched across the world. Those files are identical for every user, yet every request travels all the way to the app servers.
Fix. Serve static files from a CDN (content delivery network). A CDN is a network of cache servers, called edge servers or points of presence, spread across many cities. When a user in Toronto requests app.js, DNS directs them to a nearby edge. If the edge has a copy, it returns it immediately. If not, it fetches the file once from Recipebox's servers (the origin), keeps it and serves the next thousand Toronto users from local storage.
Toronto user --> [CDN edge Toronto] --(miss only)--+
Paris user --> [CDN edge Paris] --(miss only)--+
Pune user --> [CDN edge Mumbai] --(miss only)--+
v
+-------------------+
| Origin: LB + apps |
+-------------------+
The numbers explain why this matters. If 1 million users each load pages 5 times a day, and each page pulls 2 MB of static files (ignoring browser caching), that is 10 TB per day leaving your servers. With a CDN and a high hit ratio, almost all of it is served from the edge instead, and the files travel a much shorter distance to each user.
New trade-offs:
- Stale files after a deploy. If you publish a new
app.jsbut the edge holds the old one, users get a broken mix. The standard fix is to put a content hash in the filename (app.3f9a1c.js): a new version gets a new name, so the old cached copy is simply never requested again. - Cost and another vendor. CDNs charge for data transfer and add configuration to manage.
- Not everything is cacheable. Personal pages (your saved recipes) must never be cached in a shared edge, or one user may see another's data.
Networking and edge delivery covers cache keys, push vs pull CDNs, and invalidation.
Stage 7: Object storage for photos
Symptom. Photo uploads are growing fast. About 2% of the 1 million early users uploaded a 3 MB photo each day: roughly 60 GB per day, about 22 TB per year, and it has grown since. Photos sit on app servers' local disks, which causes three problems. Disks fill up. A photo uploaded to server 1 is missing on server 2, which breaks the "stateless" rule. And if a server dies, its photos are lost.
Fix. Store photos in object storage: a managed service that stores files ("objects") of any size under a key, such as photos/2026/10/recipe-42.jpg, and serves them over HTTP. Object storage is designed to hold enormous amounts of data cheaply and durably by keeping several copies across machines and often across buildings. It does not support editing part of a file or running queries; you put a whole object, get it back or delete it.
The upload flow becomes:
- The app server checks the user is allowed to upload and creates a database row for the photo.
- The client uploads the file to object storage, often directly using a short-lived upload URL that the app server signs, so the large file never passes through the app servers.
- The database stores only the object's key, not the bytes.
- The CDN uses object storage as its origin for photos, so readers get photos from a nearby edge.
+-----+ static + photos +-----------------+
users --------->| CDN |------------------>| Object storage |
| +-----+ (on miss) | (photos, files) |
| +-----------------+
| API +----+ +-----------+ ^ upload
+------------>| LB |-->|App servers|----------+ (signed URL)
+----+ +-----------+
| \
+-------+ +---------+ +---------+
| Cache | | Primary |-->| Replicas|
+-------+ +---------+ +---------+
New trade-offs:
- Two places to keep in sync. The database row and the object can disagree. If the upload fails after the row is created, you have a recipe pointing to a missing photo. You need a clean-up job or a "pending" status that becomes "ready" only after the upload succeeds.
- Access control. Public recipe photos can be public objects. Private drafts need signed, expiring URLs.
- Latency for the first byte is usually higher than reading a local disk, which is why the CDN sits in front.
Stage 8: Asynchronous work with queues
Symptom. Posting a recipe now feels slow, about 4 seconds. Engineers trace the request and find the app server does everything while the user waits:
- Save the recipe.
- Generate three thumbnail sizes from the photo.
- Update the search index.
- Email every follower of the author.
Steps 2 to 4 are slow, and when the email provider has an outage, posting fails entirely, even though the recipe itself saved fine. Consider just the email step: sending 1 million notification emails a day at 200 ms each would take about 55 hours if done one at a time, which is more than a day's worth of work.
Fix. Split the work into "must happen now" and "can happen shortly". The app server saves the recipe, puts a small message on a message queue saying "recipe 42 was created", and immediately responds to the user. A message queue is a component that stores messages durably until a consumer processes them. Separate worker processes read messages from the queue and do the slow work, retrying if something fails.
+-----------+ 1. save +---------+
|App servers|------------>| Primary |
+-----------+ +---------+
| 2. publish "recipe 42 created"
v
+-----------------+ +-----------------------------+
| Message queue |---->| Workers |
+-----------------+ | - make thumbnails |
| - update search index |
| - send follower emails |
+-----------------------------+
The user now waits only for step 1, perhaps 100 ms. If the email provider is down, messages wait in the queue and are retried later. If uploads spike, the queue absorbs the burst and workers catch up, which is called load levelling. You can add workers independently of app servers.
New trade-offs:
- Eventual completion. The thumbnail might appear a few seconds after the recipe. The UI must handle "still processing".
- Duplicates and ordering. Most queues deliver a message at least once, so a worker may see the same message twice. Workers must be idempotent, meaning running the same job twice has the same effect as running it once (for example, "send email for comment 9001 unless already sent").
- Back pressure. If messages arrive faster than workers can handle, the queue grows without bound. You need limits and alerts, and sometimes you must reject new work with a "try again later" response.
- Harder debugging. A single user action is now spread across processes and time. You need request IDs and tracing.
See message queues and Kafka and event streaming for delivery guarantees and partitioning.
Stage 9: Splitting the database (federation and sharding)
Symptom. Recipebox reaches 50 million DAU. Reads are fine: caches and replicas absorb them. But writes keep growing: comments, likes, saves and new recipes. The primary database handles every write and its disk and CPU are saturated. The tables are also so large (billions of rows of likes) that index maintenance and backups take hours. Vertical scaling of the primary is already maxed out.
There are two main ways to split a database, and they are often used in sequence.
Step 9a: Federation (split by function)
Federation, also called functional partitioning, puts different kinds of data in different databases. Recipebox splits into:
- a users database (accounts, profiles),
- a recipes database (recipes, ingredients, photos' keys),
- an engagement database (comments, likes, saves).
+-----------+
|App servers|
+-----------+
/ | \
+---------+ +---------+ +-------------+
| Users DB| |Recipes | |Engagement DB|
|(+replicas) |DB (+rep)| | (+replicas) |
+---------+ +---------+ +-------------+
Each database now gets only a fraction of the writes, and each can be tuned separately. Trade-off: you can no longer join users to comments in one SQL query. The app must query two databases and combine results, and a transaction that spans two databases is much harder.
Step 9b: Sharding (split by key)
Federation runs out when a single kind of data is still too big. Likes alone outgrow one engagement database. Sharding (horizontal partitioning) splits one table's rows across many databases, called shards, using a shard key. For example, store each like on shard number hash(recipe_id) mod 8. All likes for recipe 42 live on one shard, and the eight shards each carry roughly one eighth of the writes.
+--------------------------+
App server ---->| Router: shard = hash(id) |
+--------------------------+
/ / | \ \
+-------+ +-------+ +-------+ +-------+ ...
|Shard 0| |Shard 1| |Shard 2| |Shard 3| (each with replicas)
+-------+ +-------+ +-------+ +-------+
New trade-offs (sharding is the most expensive step so far):
- Choosing the key is hard to undo. A good key spreads load evenly and keeps data that is read together on the same shard. A bad key creates a hot shard: if one celebrity recipe gets 30% of all likes, its shard melts while others sit idle.
- Cross-shard queries are slow. "Top 10 most-liked recipes this week" must ask every shard and merge the results.
- Resharding. With
hash mod 8, going to 9 shards moves most keys. Consistent hashing (explained in the building blocks lesson and scalability) reduces that movement. - Operational cost. Backups, schema changes and monitoring now happen on many databases.
At this stage teams also denormalise: store a recipe's like_count directly on the recipe row instead of counting likes across shards on every page view. That makes reads cheap at the price of keeping the copy updated.
Interview tip
Do not shard first. Say the order out loud: "I would scale up, add read replicas and caching, split by feature, and only shard the tables that still do not fit, choosing the shard key from the most common query." That sentence shows judgment.
Stage 10: Multiple regions
Symptom. Two problems arrive together. Users in Europe and North America see about 150 to 250 ms of extra delay on every API call, because the speed of light across continents cannot be cached away for personalised requests. Then a power incident takes the whole Mumbai data centre offline for two hours, and Recipebox is down everywhere.
Fix. Run the system in more than one region (a region is a geographic area containing one or more data centres). Each region gets its own load balancers, app servers, caches and database copies. A global routing layer, often DNS-based or an anycast network, sends each user to the nearest healthy region.
+-----------------------+
users ---------> | Global DNS / routing |
+-----------------------+
/ \
+------------------+ +------------------+
| Region: Mumbai | | Region: Frankfurt|
| LB, apps, cache | | LB, apps, cache |
| DB (primary for |<--->| DB (primary for |
| Indian users) | rep | EU users) |
+------------------+ +------------------+
\ /
+------------------------------+
| Object storage (replicated) |
| + CDN at the edge |
+------------------------------+
There are two broad options:
| Option | How it works | Good | Painful |
|---|---|---|---|
| Active-passive | One region serves traffic; another is a warm standby that takes over on failure | Simple data model, one writer | Standby costs money while idle; failover takes minutes; no latency win for distant users |
| Active-active | Every region serves traffic at once | Low latency everywhere; survives a region loss | Writes in two regions can conflict; much harder data design |
For active-active, a common approach is to give each user a home region that owns their writes, so most writes never conflict, while reads of public content are served from local replicas.
New trade-offs: cross-region replication is slow (tens to hundreds of milliseconds), so you must decide which data must be strongly consistent everywhere (account balances, usernames) and which can be eventually consistent (like counts, feeds). Costs roughly double. Testing failover becomes a regular drill, because a failover you have never practised is likely to fail. See multi-region design for the full treatment.
The whole journey on one page
Here is the final Recipebox architecture in one region, with the stage that added each piece:
users
|
v
[DNS / global routing] ............................ stage 10
| \
v v
[CDN edge] ----> [Object storage] .............. stages 6, 7
|
v
[Load balancer pair] ........................... stage 3
|
v
[Stateless app servers x N] .................... stage 3
| | \
v v v
[Cache] [Message queue]--[Workers] ............ stages 5, 8
| |
v v
[Federated + sharded databases, each with
a primary and read replicas] .................. stages 1, 4, 9
+
[Search index, fed by workers] ................. stage 8
And here is the same story as a table you can revise from:
| Stage | Symptom | Fix | New trade-off |
|---|---|---|---|
| 0 | (none yet) | One server | Single point of failure |
| 1 | App and DB fight for resources | Separate DB server | Network hop per query |
| 2 | CPU maxed at peak | Bigger machine | Ceiling, cost, still a SPOF |
| 3 | One app server not enough; outages | Load balancer + stateless servers | Shared session store; LB redundancy |
| 4 | DB overloaded by reads | Read replicas | Replication lag; writes do not scale |
| 5 | Same reads repeated | Cache (cache-aside) | Stale data, invalidation, cold start |
| 6 | Slow static files far away | CDN | Versioning, never cache private pages |
| 7 | Photos fill disks, break statelessness | Object storage | DB row and object can disagree |
| 8 | Slow requests, fragile side effects | Queue + workers | Duplicates, eventual completion, back pressure |
| 9 | Writes overload the primary | Federation, then sharding | No joins, hot shards, resharding |
| 10 | Distant users slow; region outage | Multiple regions | Conflicts, cost, failover drills |
Things that grow alongside every stage
Three supporting pieces are not tied to one stage. They become more important at each step:
- Monitoring and alerting. You can only find the next bottleneck if you measure it: request rate, error rate, latency percentiles, CPU, queue depth, replication lag, cache hit ratio. See observability and performance.
- Automation. By stage 3 you cannot hand-configure servers. Deploys, scaling and failover need scripts or managed services.
- Service boundaries. As teams grow, the single app is often split into services (for example, a search service and a notification service). That is an organisational decision as much as a technical one; see microservices.
How to use this story in an interview
Interviewers rarely say "grow this from one server". But almost every design question becomes easier when you walk this ladder in your head:
- Estimate the traffic (the next-but-one lesson shows how).
- Ask: does one server plus one database handle it? For many internal tools, it does.
- If not, which resource breaks first? Reads, writes, storage, bandwidth or latency for distant users?
- Apply the matching fix from the table, and state its trade-off.
- Stop when the requirements are met. Extra boxes without a reason cost points.
Common mistake
Drawing the stage-10 diagram for every problem. A design for 10,000 internal users that includes sharding and three regions tells the interviewer you are reciting, not reasoning. Match the architecture to the numbers.
Interview tip
When you add a component, use a three-part sentence: "Because of symptom, I add component, which costs us trade-off, and I handle that by mitigation." For example: "Because 90% of traffic is reads, I add read replicas; that introduces replication lag, so a user's own reads go to the primary for a few seconds after they write."
Interview questions
Q1. Why would you start with a single server instead of a scalable architecture?
Because at small scale a single server is cheaper, faster to build and easier to debug, and its limits are far away. Complexity such as load balancers, caches and shards adds failure modes and operational work. The right approach is to know which component breaks first and add the fix when the numbers say so.
Q2. What is the difference between vertical and horizontal scaling?
Vertical scaling makes one machine bigger (more CPU, RAM, faster disks); horizontal scaling adds more machines and spreads work across them. Vertical is simpler and needs no code changes, but it has a hard ceiling, rising cost and remains a single point of failure. Horizontal scaling has no practical ceiling and tolerates machine failures, but it requires stateless services or data partitioning.
Q3. What does "stateless app server" mean and why does it matter?
A stateless server keeps nothing about a user between requests in its own memory or disk; sessions, files and data live in shared stores. This makes servers interchangeable, so the load balancer can send any request to any server, and you can add or remove servers freely. Without it, users get logged out or lose uploads when traffic moves between servers.
Q4. Read replicas improved read throughput. Why did they not fix write throughput?
Every replica must apply every write, and all writes still go through the single primary. Adding replicas multiplies the machines that can answer reads but does nothing to reduce the primary's write load. To scale writes you must split the data, by federation or sharding, so different writes go to different primaries.
Q5. A user posts a comment, refreshes and does not see it. What happened and how do you fix it?
Most likely the refresh read from a replica that had not yet applied the write, which is replication lag. Fixes include reading from the primary for that user for a short window after a write (read-your-own-writes), reading from the primary for the specific item just written, or showing the new comment from the client's local state. The general lesson is that replicas are eventually consistent.
Q6. What is cache-aside, and what are its risks?
In cache-aside, the app checks the cache first; on a miss it reads the database, stores the result in the cache with a TTL and returns it. Risks are stale data after updates, a burst of misses when a hot key expires or the cache restarts, and memory limits that force eviction. Mitigations include invalidating on write, short TTLs, and letting only one request rebuild a hot key.
Q7. Why put user photos in object storage instead of the database or app server disks?
App server disks break statelessness, fill up and lose data when the server dies. Databases are optimised for small structured rows and become expensive and slow with large binary blobs. Object storage is built for large files: cheap, durable through replication, served over HTTP and easy to place behind a CDN; the database stores only the object's key.
Q8. When should work be moved to a message queue?
When it is slow, can happen a little later, or depends on an unreliable external service, such as sending emails, resizing images or updating a search index. A queue lets the user get a fast response, absorbs traffic bursts and lets failed jobs be retried. The cost is eventual completion, possible duplicate deliveries that require idempotent workers, and more complex debugging.
Q9. Federation vs sharding: what is the difference?
Federation splits a database by function, so users, recipes and comments live in separate databases. Sharding splits one table's rows across many databases by a key, such as hash(user_id). Federation is simpler and often comes first; sharding is needed when a single kind of data is still too large or write-heavy for one machine.
Q10. How do you choose a shard key?
Pick a key that spreads data and traffic evenly and that matches your most common query, so most requests touch one shard. For a chat app, conversation_id keeps a conversation's messages together. Avoid keys with skew, such as country or a celebrity user, and avoid monotonically increasing keys with range sharding, which send all new writes to the last shard.
Q11. What problems does a CDN solve, and what can go wrong?
A CDN serves cacheable content from edge servers close to users, reducing latency and offloading bandwidth from your origin. Problems include serving stale files after a deploy (solved with content-hashed filenames), caching private responses by mistake, and origin overload when the cache is purged. It mainly helps static or public content, not personalised API calls.
Q12. Active-passive vs active-active multi-region?
Active-passive serves from one region and fails over to a standby; it is simpler because there is one writer but wastes standby capacity and does not reduce latency for distant users. Active-active serves from all regions at once, giving lower latency and better resilience, but concurrent writes in different regions can conflict. A common middle ground is assigning each user a home region for writes.
Q13. Your system has a load balancer. Is it still a single point of failure?
It can be. A single load balancer instance is itself a SPOF, so production systems use a redundant pair with failover, multiple instances behind DNS, or a managed load balancer that is already distributed. The same question applies to every component: always ask what happens when this box dies.
Key takeaways
- Every component in a large architecture exists to solve a specific symptom; learn the symptom and you can rebuild the diagram.
- Start simple. One server, then a separate database, is correct for small traffic.
- Vertical scaling is easy but has a ceiling and stays a single point of failure; horizontal scaling needs stateless services.
- Read replicas and caches scale reads; only federation and sharding scale writes.
- CDNs and object storage move heavy, static bytes away from your app servers and closer to users.
- Queues make slow or unreliable work asynchronous, at the cost of duplicates and eventual completion.
- Each fix introduces a trade-off. Naming it, and a mitigation, is what interviewers listen for.
- Match the architecture to the numbers; stop adding boxes when requirements are met.
Next lesson
Continue with System Design Fundamentals.

