
System design is the big-picture plan for how servers, databases, caches and queues work together when lots of people use an app at once. These system design basics explain how large apps stay fast and keep running when traffic suddenly floods in.
This guide follows the same path as our video. It starts with a single server and adds one idea at a time: scaling, load balancing, replication, sharding, caching, CDNs, the CAP theorem and message queues. Each step explains what the piece does, why you need it, and what it costs you.
- System design is planning how all the pieces of an app cooperate at scale, not writing one function.
- Start with requirements: features, scale and how fast things must feel.
- Scale out with many servers behind a load balancer, then fix the database with replication (reads) and sharding (writes).
- Caches and CDNs keep data close, so most requests never touch the database or cross the planet.
- The CAP theorem forces a choice between consistency and availability during network splits, and queues keep failures contained.
- 0:00Intro
- 0:25When the traffic hits, everything breaks
- 1:16What is system design?
- 1:50Start with requirements
- 2:55Vertical vs horizontal scaling
- 3:54Meet the load balancer
- 5:00Your database is the bottleneck
- 5:48Replication: copies that share the load
- 6:48Sharding: split the data
- 7:50Caching changes everything
- 8:51CDNs bring data next door
- 9:47The CAP theorem
- 10:39Bank or feed?
- 11:32Message queues absorb the chaos
What is system design and why does it matter?
Apps rarely die when nobody uses them. They die the day everybody suddenly does. Imagine launching a shop and having a post go viral. A wave of visitors arrives, pages spin forever and checkout times out. Those visitors don't wait around. They close the tab, and most never come back.
Downtime costs real money. Orders vanish, support inboxes fill up, and the story online becomes "that app that crashed" instead of "that cool new app". The apps that survive aren't lucky. Someone asked early on what would happen if ten times more people showed up. Asking that question is system design.
Think of it as planning a city instead of a single house. A city needs roads, water, power and emergency services. Each piece is simple on its own. The value comes from how the pieces are arranged on purpose, so they keep working together under pressure.
- Success is a stress test you can't schedule.
- System design covers servers, databases, caches and queues working together.
- Good design is made before the traffic arrives.
How to start a system design: requirements first
Every good design begins with questions, not servers. Before drawing a single box, pin down three things: what the system must do, how big it must get, and how quickly it must respond. Skip this step and everything later rests on a guess.
Functional requirements are the features users touch. For a chat app, that means sending messages, seeing who's online and loading history. Scale is the second question. A design for a hundred people looks completely different from one for a hundred million. Latency is the third. A chat message should appear almost instantly, but a monthly report can take a while.
Getting requirements wrong gives you the right design for the wrong problem. Knowing what must be fast tells you where to spend your effort, and writing features down plainly makes it clear what you are, and are not, designing for.
- Functional needs: what the system must do
- Scale targets: users, requests and data, today and with growth
- Latency goals: how fast each part must feel
Vertical vs horizontal scaling: what's the difference?
When one server can't keep up, you have two options: make it bigger or add more of them. Vertical scaling, or scaling up, means adding CPU and RAM to one machine. It's simple and needs no code changes, but it hits a hard ceiling. Eventually there is no bigger server to buy, and top-end machines are very expensive.
Horizontal scaling, or scaling out, means adding more ordinary machines side by side. Picture the difference between buying a bigger truck and buying a fleet. A fleet has more moving parts, but you can keep adding trucks, and if one breaks down, the others keep delivering.
That resilience is why almost every big app runs on fleets, not giants. A single large machine is also a single point of failure. Scaling out solves that, but it creates a new question: something now has to split the incoming traffic between all those servers.
What does a load balancer do?
A load balancer is the traffic cop standing in front of your servers. Users never talk to the servers directly. They hit one address, and the load balancer passes each request to a server that can handle it.
The simplest strategy is round robin. The first request goes to server A, the next to B, then C, then back to A. It works like a restaurant host seating guests evenly across waiters so nobody gets swamped. Smarter balancers send work to whichever server is least busy at that moment.
A load balancer also runs health checks. If a server crashes, it stops sending traffic there, and users don't notice. Finally, it makes growth painless. You can add servers behind it at any time, and users keep using the same address.
- Spreads load so no single server melts
- Skips failed servers using health checks
- Lets you add capacity without changing the public address
Why the database becomes the bottleneck, and how replication helps
Web servers are easy to copy because they usually don't hold anything important. The data is the hard part, and every server reads and writes the same database. Imagine fifty cashiers all running to one filing cabinet. The cashiers aren't the problem; the cabinet is. Databases also keep data on disk, which is far slower than memory, so waits pile up under load.
Most apps read far more than they write. You scroll through hundreds of posts for every one you publish. Replication handles that by keeping copies of the database on several machines. One primary takes the writes and sends each change to its replicas, and reads can go to any replica.
Need more read capacity? Add another replica, like photocopying a popular library book. Replicas double as backups, since one can be promoted if the primary fails. The catch is a slight delay: right after you post, a replica might not have the change yet, so a read can be briefly stale.
- Flow: write → primary → replicas
- More replicas means more readers served at once
- Copies can lag a moment behind the primary
What is database sharding?
Replication helps reads, but every write still lands on one primary. When writes pile up, copies won't save you. You need to split the data itself. Sharding cuts data into pieces and puts each piece on its own database, so each shard handles only its slice of the writes.
When a request arrives for a user, the system looks at a shard key, such as the user ID, to decide which shard owns that user. The write goes straight to that shard and nowhere else. It's like an old phone book split into volumes: names A to H in one, I to P in another. You know exactly which volume to open.
Adding shards adds room for writes, which is how giant apps keep up with huge amounts of new data. Choose the key carefully, though. Shard by country, for example, and one very large country could overload a single shard while the others sit nearly empty.
How caching and CDNs make apps faster
Caching asks a different question: what if most requests never reached the database at all? A cache keeps frequently used data in memory, which is far faster than disk. The classic pattern checks Redis, a fast in-memory store, first. If the data is there, you're done. If not, fetch it from the database, save a copy in Redis and set it to expire after a minute.
A cache hit skips the database entirely, so it stays calm during a rush. It's like keeping your favourite mug on your desk instead of fetching it from the basement. The risk is freshness: if the real data changes, the cache might still serve the old copy, which is why expiry times matter.
A CDN, or content delivery network, applies the same idea worldwide. Edge servers hold copies of images, videos and files close to users, so someone in Tokyo gets files from a Tokyo edge instead of a server across an ocean. Shorter trips make pages faster, and the CDN absorbs most file requests, shielding your own servers when content goes viral.
- Cache hit: instant answer, no database trip
- Set expiry times to limit stale data
- CDNs cut distance and protect your origin servers
# cache-aside with Redis
let user = await redis.get(key)
if (!user) {
user = await db.findUser(id)
await redis.set(key, user, 'EX', 60)
}
The CAP theorem explained: consistency or availability?
Once data lives on many machines, the CAP theorem applies. CAP stands for consistency, availability and partition tolerance. The key moment is a network partition, when machines lose contact with each other. During that split, you must choose: refuse some requests to stay correct, or keep answering and risk being out of date.
Consistency means every user sees the same latest data, even if some requests must wait or fail. Availability means every request gets an answer, even if it's slightly stale. The right choice depends on what a wrong answer costs. A bank chooses consistency, because withdrawing the same money twice would be a disaster. A social feed chooses availability, because a blank feed drives people away while a late like does no harm.
You don't have to pick once for the whole app. A shopping site can be strict about payments and relaxed about reviews. The trade-off only bites during a partition; when the network is healthy, well-built systems usually give you both. Know which way your database leans, and don't let a default setting decide for you.
How do message queues work?
If every service calls the next one directly and waits, one slow piece can drag everything down. With a message queue, a service drops a message into the queue, like posting a letter, and gets straight back to work. Other services pick the message up and process it at their own pace.
Picture an online order. Checkout publishes one "order placed" message into Kafka. The email, inventory and analytics services each read it and do their own job. Checkout doesn't need to know who's listening, so adding a rewards service later just means having it read from the queue.
Queues also absorb spikes. During a flash sale, orders can arrive faster than emails are sent, and they simply wait in line while workers catch up. Failures stay contained too: if the email service breaks, orders still succeed, and the messages wait until it recovers.
- Decoupled: producers don't know who consumes
- Absorbs spikes instead of crushing workers
- Stops one failure from cascading through the system
Key takeaways
- System design is about arranging simple pieces on purpose, before the traffic arrives.
- Always start with features, scale and latency requirements.
- Scale out behind a load balancer; use replication for reads and sharding for writes.
- Caches and CDNs keep hot data close, but watch for stale copies.
- Choose consistency or availability per feature based on what a wrong answer costs.
- Message queues decouple services, smooth out spikes and contain failures.
Frequently asked questions
What is system design in simple terms?
System design is the big-picture plan for how servers, databases, caches and queues work together when many people use an app. It's less about the code inside one function and more about how all the pieces connect.
Is horizontal scaling better than vertical scaling?
For large systems, usually yes. Vertical scaling is simpler but hits a ceiling and leaves you with one point of failure, while horizontal scaling keeps growing as you add machines and survives individual server failures.
What is the difference between replication and sharding?
Replication copies the whole database to several machines so more reads can be served, with one primary handling writes. Sharding splits the data into pieces on separate databases so writes are spread out.
What is the difference between a cache and a CDN?
A cache keeps frequently used data in fast memory inside your system so requests can skip the database. A CDN is a worldwide network of cache servers that stores files close to users to cut travel distance.
What does the CAP theorem say?
When a network partition splits your machines, a distributed system must choose between consistency (everyone sees the latest data) and availability (every request gets an answer). When the network is healthy, well-built systems can usually provide both.
Why use a message queue instead of calling services directly?
A queue lets a service hand off work and move on instead of waiting. It absorbs traffic spikes and keeps one failing service from breaking the rest of the system.