System Design Foundations
Core concepts, numbers, and patterns to internalize before tackling individual problems. This file explains the why behind every concept, not just the what.
The 30-Minute Interview Framework
Before you know any distributed systems concept, know this framework cold. It structures how you present any design.
0–3 min → Clarify requirements
Ask: Who uses this? What are the core actions? What scale?
Produce: Functional requirements (features) + Non-functional (scale, latency, availability)
3–5 min → Capacity estimation
Compute: QPS (reads + writes), storage, bandwidth
Rule of thumb: 1M req/day ≈ 12 QPS; 1B req/day ≈ 11,500 QPS
5–10 min → High-level design
Draw the main components: clients, API gateway, services, DB, cache
Describe the happy path end-to-end
10–20 min → Deep dive on the hard problems
Pick 2–3 interesting challenges and go deep
Example: "The hard part is seat reservation under concurrency — let me explain"
20–25 min → Security, failure modes, and trade-offs
Always mention: TLS, AuthN/AuthZ, rate limiting, audit logging
Failure: what happens if the DB goes down? If the cache is cold?
25–30 min → Wrap up and questionsThe single most important habitstate your assumptions out loud. "I'm assuming 100M daily active users, 90/10 read/write ratio — does that match what you're thinking?" This shows you reason systematically, not just name-drop technologies.
Numbers Every Engineer Should Know
These numbers let you do back-of-envelope math in your head. You don't need them exact — within an order of magnitude is fine.
Latency Reference Points
| Operation | Latency | Human intuition |
|---|---|---|
| L1 cache read | ~0.5 ns | 0.5 billionths of a second |
| L2 cache read | ~7 ns | — |
| RAM read | ~100 ns | 200× slower than L1 cache |
| SSD random read | ~100 µs | 1,000× slower than RAM |
| HDD seek | ~10 ms | 100,000× slower than RAM |
| Same datacenter network round-trip | ~0.5 ms | 5× slower than SSD |
| Cross-region (US → Europe) | ~80 ms | Speed of light across an ocean |
| DNS lookup (cached) | ~1 ms | — |
What this means in practiceif a user action requires reading from disk (HDD) instead of cache (RAM), it's 100,000× slower. This is why caching is not optional — it's the difference between a 100 µs response and a 10 ms one.
System Throughput
| System | Throughput | Notes |
|---|---|---|
| Single PostgreSQL | 10k–50k reads/s, 5k–15k writes/s | Depends on query complexity and hardware |
| Redis | 100k–500k ops/s | In-memory, single-threaded, predictably fast |
| Kafka | ~500k msg/s per partition | Can scale to millions with multiple partitions |
| CDN edge node | ~100 Gbps bandwidth | Handles static assets for millions of users |
| Modern HTTP server | 200k–500k QPS (static) | For simple responses; drops sharply with DB queries |
Size Reference Points
| Item | Size | Implication |
|---|---|---|
| Single tweet | ~280 bytes | 1B tweets ≈ 280 GB — fits in one server |
| Photo (compressed) | ~200 KB | 1M photos ≈ 200 GB — needs dedicated storage |
| 1 hour 1080p video | ~4 GB | Streaming requires CDN + chunked delivery |
| 1B users × 100 bytes each | ~100 GB | Profile table fits in one machine |
| 1B users × 1 MB each | ~1 PB | Requires distributed storage |
Quick Estimation Math
1 million requests/day ÷ 86,400 sec/day ≈ 12 QPS
1 billion requests/day ÷ 86,400 sec/day ≈ 11,500 QPS
Twitter example:
500M tweets/day at 280 bytes each
= 500M × 280 bytes = 140 GB of raw tweet data per day
= 51 TB/year — needs distributed storage + archival strategy
YouTube example:
500 hours of video uploaded per minute
At 4 GB/hour = 2,000 GB = 2 TB uploaded per minute
Need parallel upload pipelines and massive object storageCAP Theorem — The Fundamental Trade-off
The concept in one sentencea distributed system that experiences a network partition (some nodes can't talk to each other) must choose between staying available (accepting requests) or staying consistent (refusing requests until the partition heals).
Why Partition Tolerance Is Non-Negotiable
Imagine you have two database servers, Server A in New York and Server B in London, and the network cable between them is cut. Requests are still arriving. What do you do?
- You cannot pretend the partition isn't happening — if you do, Server A and Server B will accept conflicting writes and diverge.
- So you must choose: do you keep serving requests (available but potentially inconsistent) or stop serving until the link is restored (consistent but unavailable)?
This is the CAP choice. Every distributed system faces this.
CP — Choose Consistency over Availability
When the network link breaks: Server B refuses all writes.
Users get errors. Bank balances stay correct.
Use this for: banking transactions, inventory counts, seat reservations,
anything where two conflicting answers is worse than no answer.
Examples: HBase, Zookeeper, etcd
AP — Choose Availability over Consistency
When the link breaks: both servers keep accepting writes.
Users get responses. Some users see stale data for a while.
Use this for: social media feeds, DNS, caches, shopping recommendations,
anything where "slightly wrong" is better than "unavailable".
Examples: DynamoDB, Cassandra, CouchDB, DNSThe rule of thumb for interviewsask "what's worse — a stale answer or no answer?"
- If stale answer is worse (money, inventory, reservations) → CP
- If no answer is worse (social content, search, recommendations) → AP
Consistency Models
Not all "consistent" is the same. These are the levels from weakest (most scalable) to strongest (most correct).
The Spectrum
WEAKEST ←────────────────────────────────────────────────────────────→ STRONGEST Eventual Monotonic Read Read-Your-Writes Session Bounded Staleness Linearizable
Eventual Consistency — the weakest. Eventually all replicas will agree, but you don't know when.
Real-world example: You like a Facebook post. Your friend in another city might not see the like count increase for a second or two. That's fine — nobody cares if the count is off by 1 for 200ms.
Monotonic Read — you'll never see an older version than the one you already read.
Real-world example: You refresh your Twitter timeline and see 200 tweets. If you refresh again, you won't suddenly see 190 tweets (as if the newer 10 were erased). You might not see new tweets yet, but you won't go backward.
Read-Your-Writes — you always see the effects of your own writes, even if others don't yet.
Real-world example: You update your profile photo on LinkedIn. You immediately see your new photo. Other users might still see the old one for a few seconds — but you always see your own change instantly.
Session Consistency — within a single user session, reads are consistent.
Real-world example: Your shopping cart. Within one browser session, the items you added are always visible. If you open a new browser (new session), there might be a brief sync delay.
Bounded Staleness — data is at most N seconds (or N versions) behind.
Real-world example: A news site CDN caches articles for 30 seconds. You might see an article that's up to 30 seconds old, but never older. Used in Azure Cosmos DB's strong-bounded mode.
Linearizability (Strong Consistency) — every read reflects the most recent write. Behaves as if there's only one copy of the data.
Real-world example: A bank account balance. Every ATM withdrawal must see the exact current balance — not a cached version from 200ms ago. If two ATMs process withdrawals simultaneously on a $100 account, exactly one must see $100 and one must see $0 (or fail). No double-spend. Used by: Google Spanner, single-node PostgreSQL, etcd (for cluster coordination).
Database Selection Guide
The most common interview mistake: reaching for PostgreSQL for everything. Different storage problems need different storage engines.
The Decision Tree
Do you need transactions across multiple tables?
Yes → Relational (PostgreSQL, MySQL)
Is the data write-heavy with simple access patterns (no complex queries)?
Yes → Wide-column (Cassandra, HBase) or Time-series (InfluxDB, TimescaleDB)
Do you need full-text search?
Yes → Search engine (Elasticsearch) as a secondary store, not primary DB
Is this a cache/session store?
Yes → Key-value in-memory (Redis, Memcached)
Are you storing files, images, video?
Yes → Object storage (S3, GCS)
Do you have a deeply connected graph (social network, fraud detection)?
Yes → Graph DB (Neo4j, Amazon Neptune)
Otherwise → start with PostgreSQL and shard/switch laterReference Table
| Type | Pick when… | Real-world examples | Trade-off to mention |
|---|---|---|---|
| Relational (PostgreSQL) | You need JOINs, transactions, complex queries, foreign keys | User accounts, orders, inventory | Hard to scale writes horizontally; shard only as last resort |
| Wide-column (Cassandra) | High write throughput, time-series data, queries are always on partition key | IoT sensor data, activity logs, Netflix watch history | No JOINs; design your partition key carefully or you get hot spots |
| Document (MongoDB) | Hierarchical / nested data, schema changes frequently | Product catalogs, CMS content, event logs | Consistency trade-offs; avoid transactions across documents |
| Key-value (Redis) | Cache, sessions, counters, leaderboards, pub/sub | Session store, rate limit counters, hot data cache | No query capability; data must fit in memory (or spill to disk) |
| Graph (Neo4j) | Data is relationships: social graph, recommendation, fraud rings | LinkedIn "people you may know", fraud detection | Not general-purpose; poor at bulk writes |
| Search (Elasticsearch) | Full-text search, log aggregation, faceted filters | Airbnb search, ELK Stack for logs | Not a source of truth — sync from primary DB; eventual consistency |
| Time-series (InfluxDB, TimescaleDB) | Metrics, telemetry, sensor readings (timestamp is the key) | Grafana metrics, IoT monitoring, financial tick data | Optimized for append; poor for random updates |
| Object Storage (S3) | Large blobs: images, video, documents, backups | Photos on Instagram, videos on YouTube | No query capability; strong consistency on read-after-write (S3) |
Caching
Why caching exists in one sentenceRAM is 1000× faster than SSD and 100,000× faster than a HDD. If you can answer a request from RAM instead of disk, do it.
The Cache Hit / Miss Cycle
Request arrives for user_id=123's profile:
CACHE HIT (fast path):
App → check Redis → key exists → return data in ~1ms
No database involved.
CACHE MISS (slow path):
App → check Redis → key not found
→ query PostgreSQL → return data in ~20ms
→ write result to Redis with TTL=300s
→ return data to user
Next request for same user → cache hitCaching Patterns
Cache-aside (most common) — the application manages the cache manually.
Read: App asks cache → miss → App asks DB → App writes to cache → returns to user
Write: App writes to DB → (optionally) deletes or updates cache entryUse this when: read-heavy workloads, you're retrofitting caching onto an existing system.
Risk: cache stampede — if 10,000 users all request the same key at the same moment it expires, all 10,000 go to the database simultaneously. Mitigation: mutex lock (first miss takes a lock, others wait), or probabilistic early expiration (randomly refresh before TTL expires to spread the load).
Write-through — every write goes to cache AND database, synchronously.
Write: App writes to cache → cache writes to DB → returns success
Read: Cache always has fresh data → simple readUse this when: writes and reads are balanced, freshness is critical. Trade-off: higher write latency (two writes instead of one).
Write-behind (write-back) — writes go to cache immediately, database is updated asynchronously.
Write: App writes to cache → returns success immediately
Background: cache flushes to DB asynchronously every N secondsUse this when: you need very low write latency and can tolerate small risk of data loss. Risk: if the cache server crashes before flushing, those writes are lost.
Cache Eviction Policies
When the cache is full, which items get removed?
| Policy | How it works | Best for |
|---|---|---|
| LRU (Least Recently Used) | Evict the item that hasn't been accessed for the longest time | General-purpose (most caches use this) |
| LFU (Least Frequently Used) | Evict the item accessed the fewest times overall | Stable hot data (popular content that never changes) |
| TTL (Time to Live) | Evict after a fixed time regardless of access | Freshness-sensitive data (prices, inventory) |
Redis defaultallkeys-lru (when memory is full, evict the least recently used key from all keys).
Cache Invalidation — The Hard Problem
"There are only two hard things in Computer Science: cache invalidation and naming things." — Phil Karlton
How do you keep the cache consistent with the database?
let the cache entry expire naturally. Simple. Risk: up to TTL window of stale data.
when the DB is written, send an event to invalidate the cache key. More complex, but more fresh. Risk: race conditions between the write and the invalidation.
include a version number in the cache key (user_123_v5). When the user's record changes, the key becomes user_123_v6 — the old key is now unreachable (eventually evicted). Simple, no explicit invalidation needed.
Scaling Patterns
Vertical vs Horizontal Scaling
| Vertical (scale up) | Horizontal (scale out) | |
|---|---|---|
| What it means | Bigger machine (more CPU, RAM) | More machines |
| Real-world analogy | Replace a car with a truck | Add more cars |
| Ceiling | Physical hardware maximum | Theoretically unlimited |
| Cost | Expensive; price jumps at high end | Linear with load |
| Failure risk | Single point of failure | Any one node failing is fine |
| When to use first | Always try this first — simpler | When vertical hits ceiling |
The interview answerstart with vertical scaling (it's simpler, no code changes needed). Switch to horizontal when you hit the hardware ceiling or need geographic redundancy.
Load Balancing
A load balancer sits in front of your servers and distributes incoming requests. Without it, one server gets everything and the others are idle.
┌─→ Server A (handles 33% of traffic)
Users → LB ────────┼─→ Server B (handles 33% of traffic)
└─→ Server C (handles 33% of traffic)Algorithms
Server A, then B, then C, then A again. Simple, works well when requests take similar time.
next request goes to whichever server has the fewest active connections. Better when requests vary in duration (some take 1ms, some take 5s).
same client always goes to the same server. Use for stateful sessions (but better to make servers stateless and store state in Redis).
server with more capacity gets more traffic. Use when servers have different specs.
L4 vs L7 load balancers
(TCP level): fast, routes by IP/port, can't inspect HTTP headers or URLs. Use for very high throughput.
(HTTP level): can route by URL path, headers, cookies. Can do SSL termination (decrypt HTTPS at LB, send plain HTTP to servers). This is what AWS ALB, nginx, and HAProxy do.
Database Scaling
Read replicasthe simplest horizontal DB scaling. All writes go to the primary. Reads spread across multiple replicas.
Writes → Primary DB
↓ replication (async)
Reads → Replica 1, Replica 2, Replica 3Trade-off: replica lag. A replica might be 50–200ms behind the primary. Users who just wrote data might read stale data from a replica. Fix: route writes and reads-after-write to the primary; route all other reads to replicas.
Shardingsplit the database into multiple independent databases (shards), each owning a portion of the data.
user_id 1–1,000,000 → Shard 1 (its own PostgreSQL instance)
user_id 1,000,001–2M → Shard 2
user_id 2,000,001–3M → Shard 3Why it's complex: JOINs across shards are expensive or impossible. Rebalancing when a shard gets too large is painful. Pick the shard key carefully — if everyone queries by user_id, shard by user_id. Don't shard until you absolutely need to.
CQRS (Command Query Responsibility Segregation)instead of one database for both reads and writes, have two separate models. Writes go to a normalized relational DB. Reads go to a denormalized "read model" optimized for queries (could be a different database type — Elasticsearch, Redis, a pre-aggregated table).
Real-world example: Twitter's tweet storage (Cassandra for writes) vs tweet search (Elasticsearch for full-text search). Two different systems, kept in sync via an event stream.
Message Queues
Why Queues Exist
Imagine a user uploads a photo to Instagram. Three things need to happen: resize the photo to multiple resolutions, run it through content moderation, and notify followers. You could do all three synchronously — but then the upload API call takes 5 seconds. Or you can:
- Accept the upload, save the original photo.
- Return success to the user immediately.
- Push a "photo uploaded" event to a queue.
- Three separate worker services consume the event and do their jobs independently.
The user sees instant success. The heavy work happens in the background, decoupled.
What queues give you:
Decoupling: Producer doesn't know or care about consumers
Spike absorption: Consumer works at its own pace; queue buffers the backlog
Retry: Failed jobs can be retried without bothering the producer
Fan-out: Multiple consumers can process the same event (photo resize + moderation)Kafka vs RabbitMQ vs SQS
| Property | Kafka | RabbitMQ | SQS (AWS) |
|---|---|---|---|
| Mental model | Append-only log | Traditional queue | Managed AWS queue |
| Consumer model | Pull — consumer tracks its own offset (position in the log) | Push — broker delivers to consumer | Poll — consumer asks for messages |
| Retention | Messages stay for days/weeks; can replay from any point | Deleted after consumed | Up to 14 days; gone after consumed |
| Ordering | Strict within a partition; none across partitions | Per-queue | Best-effort; FIFO queue option |
| Throughput | Very high (millions/s with partitioning) | Medium | High |
| Best for | Event streaming, audit trail, replay, fan-out to many consumers | Task queue (one consumer processes each job) | Simple async jobs on AWS |
When to pick Kafkaevent streaming, event sourcing, systems that need to replay events, fan-out to many independent consumers, audit logs.
When to pick RabbitMQwork queues where exactly one worker processes each job, complex routing rules, need acknowledgment guarantees.
When to pick SQSyou're on AWS and want something managed with zero ops overhead.
Delivery Guarantees
| Guarantee | What it means | Risk | How to handle |
|---|---|---|---|
| At-most-once | Message is delivered 0 or 1 times | Can lose messages | Acceptable for non-critical metrics |
| At-least-once | Message is delivered 1 or more times | Can duplicate messages | Make consumers idempotent (processing the same message twice has the same effect as once) |
| Exactly-once | Message is delivered exactly once | Hard to achieve | Kafka transactions + idempotent producers; expensive |
The practical answerdesign for at-least-once delivery and make your consumers idempotent. Use a database deduplication key (message_id → processed = true) to detect and skip duplicates.
Consistent Hashing
The Problem It Solves
Imagine you have a cache cluster with 3 nodes (A, B, C). You decide which node holds a key with node = hash(key) % 3. This works until you add a 4th node.
With 3 nodes: hash("user:123") % 3 = 1 → Node B
With 4 nodes: hash("user:123") % 4 = 3 → Node D
Almost every key now maps to a different node. Your entire cache is invalidated overnight. The database gets a sudden surge of cache misses.
Consistent hashing solves this: when you add or remove a node, only the keys that were owned by the added/removed node need to move. All other keys stay put.
How It Works
Imagine a circle (ring) representing all possible hash values from 0 to 360 degrees. You place nodes at positions around the ring:
Ring (0° → 360°):
Node A (at 90°)
/
0° ──────────── 90° ── Node A
│ │
│ Ring │
│ │
Node C (270°) ─── 180° ── Node BTo find which node holds a key: hash the key → get an angle → move clockwise → first node you hit owns the key.
"user:123" hashes to 150° → move clockwise → hits Node B at 180° → Node B owns it
"user:456" hashes to 220° → move clockwise → hits Node C at 270° → Node C owns itAdding a node (Node D at 210°):
Before: keys between 180° and 270° were owned by Node C
After: keys between 180° and 210° → Node D (only these move)
keys between 210° and 270° → still Node C (unchanged)Only ~1/N of keys move when you add/remove a node (N = number of nodes). All other keys stay on the same node.
Virtual nodesreal implementations give each physical node multiple positions on the ring (e.g., Node A appears at 90°, 200°, 340°). This ensures even key distribution even with a small number of nodes, and allows fine-grained load balancing.
Used in: Redis Cluster, Apache Cassandra, Memcached (ketama), DynamoDB.
API Design Principles
Idempotency
An operation is idempotent if calling it multiple times has the same effect as calling it once.
Idempotent: GET /users/123 (reading twice = same result)
PUT /users/123 (setting the same value twice = same state)
DELETE /users/123 (deleting twice = still deleted)
NOT idempotent: POST /orders (submitting twice = two orders!)Why it matters: networks are unreliable. If a payment request times out, should you retry? Only if the operation is idempotent. To make POST /payments idempotent: require an idempotency key header (Idempotency-Key: uuid). Store a mapping of {idempotency_key → result}. If the same key arrives again, return the stored result instead of processing again.
Pagination
Never return all records at once. For large datasets:
Offset-based (simple, but slow for large offsets):
GET /posts?page=1&limit=20
GET /posts?page=1000&limit=20 ← DB must scan 20,000 rows to find the 1,000th page
Cursor-based (recommended for large datasets):
GET /posts?limit=20
Response: { "posts": [...], "cursor": "eyJpZCI6MTIzNH0=" }
Next page: GET /posts?limit=20&cursor=eyJpZCI6MTIzNH0=
← DB directly jumps to the row after the cursor; no scanningRate Limiting
Always design the API layer with rate limits. Two main algorithms:
Token bucketeach client has a bucket of N tokens. Each request consumes 1 token. Tokens refill at a constant rate. This allows bursting up to the bucket size.
Bucket: capacity=100, refill=10/second
User sends 100 requests at once → all allowed (uses all tokens)
11th request → rejected (bucket empty)
After 1 second → 10 tokens refilled → 10 more requests allowedSliding window countertrack request count in the last N seconds using two counters (current window + previous window weighted by overlap). More accurate than a fixed window, very efficient (2 integers per client).
Communicate limits to clients: Retry-After: 30, X-RateLimit-Limit: 100, X-RateLimit-Remaining: 0.
Versioning
Always version your APIs:
URL versioning (simplest): /v1/users, /v2/users
Header versioning: Accept: application/vnd.api+json; version=2Never break an existing API version. Run old and new versions in parallel during migration.
Distributed Systems Failure Modes
These are the failure patterns you must know and explain mitigations for.
Single Point of Failure (SPOF)
Any component that, if it fails, takes the whole system down.
Example: Your application has one Redis instance for sessions. Redis goes down → all users are logged out.
MitigationRedis Sentinel (automatic failover) or Redis Cluster (sharded + replicated). Apply the same thinking to databases, message brokers, and load balancers.
Cascading Failure
Service A calls Service B. Service B is slow. Service A's threads pile up waiting for B. Service A runs out of threads. Now Service C, which calls Service A, is also stuck. The slowdown cascades.
Real-world example: Amazon Prime Day 2018 — a dependency bottleneck caused a cascade that brought down multiple services.
MitigationCircuit breaker pattern. After N consecutive failures or when error rate exceeds a threshold, the circuit "opens" — subsequent calls immediately fail fast (no waiting). After a timeout, the circuit tries again ("half-open"). If success, it closes. Libraries: Hystrix (Java), Polly (.NET), resilience4py.
Cache Stampede (Thundering Herd)
A popular cache key expires. 10,000 simultaneous requests all miss and all hit the database at once.
Mitigationmutex lock on cache miss (first miss locks, populates; others wait), or probabilistic early expiration (refresh slightly before TTL expires to spread the load over time).
Split-Brain
Two database nodes each think they are the primary. Both accept writes. Data diverges.
Example: network partition causes two PostgreSQL nodes to both believe they're the primary. After the partition heals, you have conflicting writes on both.
Mitigationquorum-based consensus (Raft, Paxos). A node can only become primary if it has acknowledgment from a majority (quorum) of nodes — impossible for two nodes to both achieve quorum simultaneously.
Hot Partition
In a sharded database or Kafka, all traffic concentrates on one shard.
Example: you shard Kafka by
user_id. Justin Bieber has 100M followers and posts constantly. Hisuser_idpartition gets 100× the traffic of any other.
Mitigationadd a random suffix to the partition key for hot users (user_id + random(0,9) → 10× spread). For Kafka, use more partitions for high-volume topics.
Security Fundamentals for System Design
Weave these naturally into your design. Don't wait to be asked.
Authentication vs Authorization
"who are you?" → verify identity → issue a token (JWT, session cookie, API key)
"what are you allowed to do?" → check permissions → RBAC, ABAC, OPA policy
JWT best practicesshort-lived access tokens (15 minutes) + long-lived refresh tokens. Store refresh tokens in HttpOnly cookies (not localStorage — localStorage is XSS-accessible). Access tokens in memory only.
The Standard Security Stack for Any System Design
Client → CDN/WAF (DDoS protection, OWASP rules)
→ API Gateway (rate limiting, auth token validation)
→ Service (AuthZ check, input validation)
→ Database (parameterized queries, column-level encryption for PII)
→ Audit log (every mutation, immutable, off to SIEM)Mention this stack in every design, even if briefly. It signals security-aware thinking.
Data Classification
Always ask: "what data is sensitive here?" Then apply the right controls:
Public data (posts, public profiles): no special encryption; CDN cacheable
Internal data (user emails, phone numbers): encrypt at rest (column-level AES-256)
Sensitive data (passwords, SSNs, payment cards): bcrypt/Argon2 for passwords;
tokenize card numbers (store token, send real number only to payment processor)
Audit data: immutable, separate storage, append-onlyAudit Logging
Every system design interview should mention audit logging for any financial, healthcare, or user-data system:
What to log: actor (user/service), action (read/write/delete), resource (table + ID),
timestamp, source IP, outcome (success/failure)
Where to store: separate append-only log database; ship to SIEM (Splunk, Chronicle)
Why separate: so a compromised application cannot delete or modify audit recordsInput Validation
At the API boundary:
- Type check, length limits, allowed character sets
- SQL injection: always use parameterized queries (
?placeholders, never string concatenation) - SSRF: if your API accepts user-supplied URLs, block requests to
169.254.169.254(cloud metadata),10.x.x.x,172.16.x.x,192.168.x.x(internal networks) - XXE: disable XML external entity processing when parsing XML
The Big Trade-offs
Memorize these. Every system design decision is one of them.
| Decision | Option A | Option B | How to choose |
|---|---|---|---|
| SQL vs NoSQL | PostgreSQL (transactions, JOINs, strong consistency) | Cassandra/DynamoDB (horizontal write scale, flexible schema) | Need transactions or complex queries → SQL. Need to scale writes to millions/s → NoSQL |
| Strong vs Eventual Consistency | Strong (linearizable — always fresh) | Eventual (AP — may be slightly stale) | Money/inventory/reservations → strong. Social content/analytics → eventual |
| Sync vs Async | REST/gRPC (returns when done, caller waits) | Queue + worker (returns immediately, work happens later) | Need immediate result (checkout) → sync. Can tolerate delay (email, transcoding) → async |
| Fan-out on Write vs Read | Push: pre-compute each follower's feed at write time (fast reads) | Pull: compute feed at read time from followed accounts (fast writes) | Most users → push. Celebrity accounts (100M followers) → hybrid (push for normal users, pull for celebs) |
| Cache-aside vs Write-through | Lazy (only cache on first read miss) | Eager (cache on every write) | Read-heavy, cache only hot data → cache-aside. Write-then-read patterns → write-through |
| UUID vs Sequential ID | UUID (globally unique, no coordination needed) | Sequential (sortable, compact, index-friendly) | Distributed ID generation → UUID. Single-node or performance-critical → sequential |
| Bloom filter vs Hash set | Bloom filter (probabilistic, memory-efficient, false positives possible) | Hash set (exact, uses more memory) | "Probably not" is good enough (spam filter, web crawler dedup) → Bloom. Need certainty → Hash set |
| Push vs Poll | Push (server sends when data arrives, WebSocket/SSE) | Poll (client asks periodically, HTTP) | Real-time chat, live updates → push. Occasional sync (email check) → poll |
| Monolith vs Microservices | Monolith (simple, one deploy, one database) | Microservices (independent scale, independent deploy, isolated failures) | Start monolith. Break into services when you have different scaling needs or team boundaries |