What Is System Design?
System design is the process of defining the architecture, components, interfaces, and data of a system to satisfy a set of requirements. Unlike coding a single feature, it deals with the whole picture: how services talk, where data lives, how the system stays up under failure, and how it scales as traffic grows from thousands to billions of requests.
There is almost never one "correct" design. Every choice is a trade-off between latency, throughput, consistency, availability, cost, and operational complexity. Good design makes those trade-offs explicit and ties them back to requirements.
Golden rule
Requirements first, numbers second, architecture last. Never draw boxes before you know the QPS, data size, and consistency needs.
Functional vs Non-Functional Requirements
Functional requirements describe what the system does — the features and behaviors. Non-functional requirements (NFRs) describe how well it does them — the qualities. NFRs are what actually drive the architecture.
| Type | Examples |
|---|---|
| Functional | "Users can post tweets", "Shorten a URL", "Upload a video" |
| Non-functional | Availability, latency, throughput, consistency, durability, scalability, security, cost |
Always pin down NFRs early: read/write ratio, expected QPS, acceptable p99 latency, durability target, and whether the workload favors consistency or availability. These numbers change every downstream decision.
Back-of-the-Envelope Estimation
Estimation converts a vague product into concrete numbers: how many servers, how much storage, how much bandwidth. The goal is order-of-magnitude accuracy, not precision. Round aggressively.
Powers of two & useful constants
2^10 = 1 Thousand = 1 KB 2^30 = 1 Billion = 1 GB
2^20 = 1 Million = 1 MB 2^40 = 1 Trillion = 1 TB
Seconds per day ~= 86,400 (~10^5, round to 100k)
Seconds per month ~= 2.5 million (2.6M for 30 days)
1 day of QPS: daily_events / 86,400
QPS from DAU
Example: Twitter-like feed
Daily active users (DAU) = 200 million
Each user reads feed 5x/day -> 1B read requests/day
Each user writes 2 posts/day -> 400M write requests/day
Average read QPS = 1,000,000,000 / 86,400 ~= 11,600 QPS
Average write QPS = 400,000,000 / 86,400 ~= 4,600 QPS
Peak factor ~2-3x average:
Peak read QPS ~= 25,000 - 35,000 QPS
Always compute peak, not just average
Traffic is bursty. Provision for peak (typically 2-3x average, higher for events) or you will fall over at the worst possible moment.
Storage & bandwidth
Storage (5 years of posts):
400M posts/day x 365 x 5 = ~730 billion posts
Each post ~= 300 bytes text + metadata
730B x 300 B ~= 220 TB (before replication)
x3 replicas ~= 660 TB
Bandwidth (write path):
4,600 write QPS x 300 B ~= 1.4 MB/s ingress
Bandwidth (read path):
11,600 read QPS x 300 B ~= 3.5 MB/s egress
(media dominates real bandwidth — 1 image = ~200 KB)
Latency Numbers Every Engineer Should Know
These "Jeff Dean numbers" (updated for modern hardware) give intuition for where time goes. The key ratios matter more than exact values — memory is ~100,000x faster than a cross-continent round trip.
| Operation | Latency | Comparison |
|---|---|---|
| L1 cache reference | ~1 ns | baseline |
| L2 cache reference | ~4 ns | 4x L1 |
| Main memory (RAM) reference | ~100 ns | 100x L1 |
| Read 1 MB sequentially from RAM | ~3 µs | memory bandwidth |
| SSD random read | ~16 µs | ~150x RAM |
| Read 1 MB sequentially from SSD | ~50 µs | NVMe |
| Round trip within same datacenter | ~0.5 ms | 500 µs |
| Disk (HDD) seek | ~5 ms | avoid on hot path |
| Round trip CA → Netherlands → CA | ~150 ms | speed of light bound |
Takeaways: cache in memory whenever possible; a cross-region round trip costs more than a million L1 accesses; and batching one network call beats a hundred small ones because the network round trip dominates.
CAP Theorem
The CAP theorem states that a distributed data store can guarantee at most two of three properties simultaneously:
| Property | Meaning |
|---|---|
| Consistency (C) | Every read sees the most recent write (linearizability) |
| Availability (A) | Every request gets a non-error response (not necessarily latest) |
| Partition tolerance (P) | System keeps working despite dropped/delayed messages between nodes |
In any real distributed system, network partitions will happen, so P is non-negotiable. The real choice is C vs A during a partition:
During a partition, a node that can't reach peers must choose:
CP -> refuse/block the request to stay consistent
(e.g. ZooKeeper, etcd, HBase, MongoDB majority)
AP -> answer with possibly-stale data to stay available
(e.g. Cassandra, DynamoDB, Riak)
PACELC — The More Complete Model
CAP only describes behavior during a partition. PACELC extends it: if Partition, choose Availability or Consistency; Else (normal operation), choose Latency or Consistency. Even with no partition, strong consistency costs latency because it requires coordination.
| System | Partition | Else |
|---|---|---|
| Cassandra / DynamoDB | PA | EL (favor latency) |
| MongoDB (default) | PA | EC (favor consistency) |
| Spanner / etcd | PC | EC |
Consistency Models
"Consistency" is a spectrum, not a boolean. Weaker models buy you lower latency and higher availability.
| Model | Guarantee |
|---|---|
| Strong (linearizable) | Reads always see the latest committed write; behaves like one machine |
| Causal | Causally related ops seen in order; concurrent ops may differ per node |
| Read-your-writes | A user always sees their own updates (even if others don't yet) |
| Eventual | If writes stop, all replicas converge — eventually. No ordering promise. |
Pick the weakest model that is still correct
A bank ledger needs strong consistency. A like counter or a "last seen" timestamp is perfectly fine with eventual consistency — and it will be far faster and more available.
Availability: The Nines
Availability is usually stated as a percentage of uptime. Each extra nine is roughly 10x harder and more expensive.
| Availability | Downtime / year | Downtime / month |
|---|---|---|
| 99% (two nines) | 3.65 days | ~7.2 hours |
| 99.9% (three nines) | 8.76 hours | ~43 min |
| 99.99% (four nines) | 52.6 min | ~4.3 min |
| 99.999% (five nines) | 5.26 min | ~26 sec |
When services are chained in series, their availabilities multiply: three 99.9% services in a request path give 0.999³ ≈ 99.7%. Redundant (parallel) components instead multiply their failure probabilities, raising availability.
SLA vs SLO vs SLI
| Term | Definition |
|---|---|
| SLI | Indicator — the measured metric (e.g. p99 latency, success rate) |
| SLO | Objective — internal target for an SLI (e.g. 99.9% of reads < 200 ms) |
| SLA | Agreement — contractual promise to customers, with penalties if breached |
SLO should be stricter than SLA so you have margin. The gap between 100% and your SLO is your error budget — spend it on releases and risk; when it runs out, freeze changes and focus on reliability.
Scaling: Vertical vs Horizontal (Intro)
Vertical (scale up) Horizontal (scale out)
Bigger machine More machines
+--------+ +--+ +--+ +--+ +--+
| 128 GB | | | | | | | | |
| 64 CPU | +--+ +--+ +--+ +--+
+--------+
simple, no code change needs LB + stateless design
hard ceiling, SPOF near-linear, fault tolerant
Vertical scaling is simplest but hits a hardware ceiling and leaves you with a single point of failure. Horizontal scaling is the foundation of large systems but requires stateless services and a load balancer.
Single Points of Failure (SPOF)
A SPOF is any component whose failure takes down the whole system: one database, one load balancer, one cache node holding critical state. Eliminate them with redundancy (replicas, multiple AZs), failover (standby promoted on failure), and no shared mutable state in the request path.
Common hidden SPOFs
A single DNS provider, one region, a shared config service, or a lone primary database. Ask "what breaks if this one box dies?" for every box in your diagram.
A Requirements-First Interview Framework (RESHADED)
A repeatable structure keeps a design interview (and real design docs) organized. RESHADED is one popular version:
| Step | What you do |
|---|---|
| R — Requirements | Clarify functional + non-functional; state assumptions |
| E — Estimation | QPS, storage, bandwidth, number of servers |
| S — Storage schema | Data model, entities, access patterns |
| H — High-level design | Major components and how requests flow |
| A — API design | Endpoints / RPCs and their contracts |
| D — Detailed design | Zoom into hard parts: sharding, caching, queues |
| E — Evaluation | Check design against requirements; discuss trade-offs |
| D — Distinctive / bottlenecks | Address SPOFs, hot spots, and unique constraints |
Interview tip
Spend the first 5 minutes on requirements and estimation. A design that ignores the numbers — or solves the wrong problem — fails no matter how elegant the boxes look.