Week 2 · System Design
A sequence of problems, decisions, and trade-offs.
The method
Begin with the simplest possible architecture. When a bottleneck appears, understand why it is happening and introduce only the component needed to solve that problem.
Requirements and reasoning should drive architecture, not the other way around.
Starting point
We only have 10–100 users.
10–100 Users
↓
Application
↓
Database
We don't need Kubernetes, multiple databases, queues, caching, or sharding.
Do not add complexity before you have a problem that requires it.
Traffic grows
100 users → 100,000 users. The architecture is still one application and one database.
First bottleneck
1,000 users → CPU 40% 10,000 users → CPU 75% 100,000 users → CPU 95–100%
Requests start waiting. Latency increases. Timeouts. Poor user experience.
Step 1
4 CPU / 16 GB RAM becomes 16 CPU / 64 GB RAM. Often the easiest first step. There is a limit.
The catch
Application Server
❌
↓
Entire application unavailable
So instead of making one machine endlessly bigger, we add more application servers.
Step 1 continued
App 1 → 2,500 requests/sec App 2 → 2,500 requests/sec App 3 → 2,500 requests/sec App 4 → 2,500 requests/sec
New problem
The user should not have to know:
Send request 1 → App 1 Send request 2 → App 3 Send request 3 → App 2
We need something in front of the application servers.
First new component
The application layer can now scale horizontally.
Problem 2
What routing algorithm should we use?
Problem 3
App 1 → /health → 200 OK App 2 → /health → 200 OK App 3 → /health → Failure
The Load Balancer temporarily removes App 3 from the backend pool. This introduces Health Checks.
Problem 4
Sticky Sessions
User A → App 1 User A → App 1 User A → App 1
A more scalable design is often:
Application Servers = Stateless
Stateless application servers are generally easier to scale horizontally.
Problem 5
Layer 4 needs IP, Port, TCP / UDP.
Client ↓ TCP Connection ↓ Layer 4 Load Balancer ↓ Backend Server
Layer 7 understands HTTP.
/api/users → User Service /api/orders → Order Service /api/chat → AI Service
Problem 6
Before: One Application Server → SPOF After: One Load Balancer → SPOF
The Load Balancer layer itself must also have redundancy, failover, and high availability.
Do not remove a Single Point of Failure by simply creating a new Single Point of Failure somewhere else.
Next bottleneck
Problem → Solution → New Problem → Deep Dive → Evolve the Architecture
Step 2A
SELECT * FROM users WHERE email = 'user@example.com';
If email is not indexed, the database may need to scan a huge amount of data.
Before scaling infrastructure, make sure the current infrastructure is being used efficiently.
Step 2B
4 CPU / 16 GB RAM / 500 GB SSD
↓
16 CPU / 64 GB RAM / Faster SSD/NVMe
Often a perfectly reasonable next step. There is still ONE DATABASE.
Step 2C
Replication can help with availability, disaster recovery, and sometimes read scalability.
Read replicas
100K Users
↓
Load Balancer
↓
Multiple App Servers
↓
┌─────────────────────┐
│ Database Layer │
├─────────────────────┤
│ Primary ← Writes │
│ Replica 1 ← Reads │
│ Replica 2 ← Reads │
│ Replica 3 ← Reads │
└─────────────────────┘
Replication maintains copies. A read replica is specifically a replica used to serve read traffic.
Step 2D
Reads → Replica 1 / Replica 2 / Replica 3 Writes → Primary
We can distribute reads. All writes are still going to ONE PRIMARY DATABASE.
Another problem
100K users → 1M → 10M → 100M users 1 billion rows
One database server has to handle all the data, all the writes, all the indexes, and all the queries.
How can we divide the data?
Step 2E
Partitioning means dividing a large dataset into smaller logical pieces.
Database
↓
┌─────────┼─────────┐
↓ ↓ ↓
Partition 1 Partition 2 Partition 3
Those partitions may still exist on ONE DATABASE SERVER.
Step 2F
Partitioning = divide the data
Sharding = distribute those partitions across machines
Shard 1 → Users 1–1M Shard 2 → Users 1M–2M Shard 3 → Users 2M–3M
Now we are scaling the database horizontally.
Shard key
hash(user_id) → shard User A → Shard 1 User B → Shard 3 User C → Shard 2
A good shard key should distribute data and traffic relatively evenly.
New problem
US → Shard 1 → 80% Europe → Shard 2 → 10% Asia → Shard 3 → 10%
We technically have three shards. But Shard 1 is still overloaded.
New problem
“Find the 100 most active users in the entire system” may need every shard.
Query → Shard 1 / Shard 2 / Shard 3 → Merge Results
Scaling solves one problem and introduces another.
One more problem
4 shards become 8 shards. Some existing data may need to move.
Moving a large amount of data while keeping the application available can be difficult, which is one reason not to jump to sharding too early.
Database evolution
1 Single Database 2 Optimize queries / indexes 3 Vertical Scale Database 4 Primary + Replicas 5 Writes → Primary, Reads → Replicas 6 Partition the data 7 Shard across multiple database servers
Even with a scalable database, why should we repeatedly hit the database for the same frequently requested data?
Cache
In general, a cache read may take around 1 ms or less, while the same request from the database may take 10–20 ms or more.
Cache path
Request → Check Cache → Cache Hit? Yes → Return No → Database → Update Cache → Return
Cache scales too
Cache 1
↗
Application → Cache 2
↘
Cache 3
This introduces partitioning, consistent hashing, replication, eviction, TTL, invalidation, and cache stampede.
Cache memory is limited
When the cache is full, which item should we remove?
Problem 5
Database → Alicia Cache → Alice
TTL — Time To Live gives the cache entry an expiration time.
Home feed → 30-second TTL Profile card → a few minutes Trending list → 60-second TTL
TTL trade-off
Short TTL ↓ Fresher data ↓ More cache misses ↓ More database traffic
Long TTL ↓ More cache hits ↓ Less database traffic ↓ Greater chance of stale data
There is no magic TTL.
Problem 6
A popular entry expires. 10,000 requests miss at once.
Only one request refreshes the cache 9,999 requests wait / use a slightly stale value
Rather than 10,000 database queries.
Another mechanism
user changes profile photo → invalidate profile:user123 permission changes → invalidate related authorization cache immediately
Good for data that must stay correct. Trade-off: fresher data, but more application complexity.
The rule
It depends on the type of data, and in many systems we use both.
TTL is easier, but explicit invalidation is safer for critical data.
The lesson
We started with one application and one database. We scaled the application, then the database, then introduced a cache. Each fix created the next question.