System design interviews and real-world architecture discussions use a lot of terminology. After going through dozens of system design problems and interviews, I put together this reference guide of the concepts that come up most often.
Scalability
The ability of a system to handle increased load by adding resources.
Vertical Scaling (Scale Up): Add more CPU, RAM, or disk to an existing machine. Simple — no code changes needed. Hard limit: the largest available machine. Single point of failure.
Horizontal Scaling (Scale Out): Add more machines to distribute load. No single ceiling. Requires the system to be stateless (or state to be externalized). Harder to implement but scales indefinitely.
Elastic Scaling: Automatically add or remove instances based on real-time load. AWS Auto Scaling, Kubernetes HPA.
Rule of thumb: Start vertical, then go horizontal when you hit limits.
Latency vs Throughput
Latency: Time for a single request to complete. Usually measured as p50, p95, p99 (percentiles). p99 = 99% of requests complete within this time.
Throughput: Number of requests handled per second (RPS). A system can have high throughput with acceptable latency, or low latency with limited throughput — the relationship is non-linear.
Little's Law: L = λW — the average number of items in a system equals the arrival rate times the average time spent. If requests arrive at 100/s and each takes 200ms, you have 20 concurrent requests in-flight.
Availability
Availability = (uptime) / (uptime + downtime)
| Availability | Annual Downtime | | ----------------- | --------------- | | 99% (2 nines) | 87.6 hours | | 99.9% (3 nines) | 8.76 hours | | 99.99% (4 nines) | 52.6 minutes | | 99.999% (5 nines) | 5.26 minutes |
Most production systems target 99.9% or 99.99%. Achieving 5 nines requires extreme engineering effort.
Availability in series: Multiple components in a chain. System fails if any component fails. A_total = A1 × A2 × A3
Availability in parallel: Redundant components. System works if any one component works. A_total = 1 - (1-A1)(1-A2)
Load Balancing
Distributes incoming traffic across multiple backend servers.
Layer 4 (Transport): Routes based on IP/port. Fast, no TLS termination. Good for high-throughput raw TCP.
Layer 7 (Application): Routes based on HTTP headers, URLs, cookies. Can do SSL termination, content-based routing, sticky sessions. Used in most web apps.
Algorithms:
- Round Robin: Each server in turn. Simple.
- Weighted Round Robin: Some servers get more traffic (for heterogeneous capacity).
- Least Connections: Route to server with fewest active connections. Good for variable-length requests.
- IP Hash: Client IP determines server. Sticky sessions without cookies.
- Consistent Hashing: Used in distributed caches — minimizes remapping when servers join/leave.
Caching
Storing a copy of data closer to the requester to reduce latency and load on origin.
Cache-Aside (Lazy Loading): Application manages the cache. On miss, load from DB and populate cache.
read(key):
value = cache.get(key)
if value is null:
value = db.get(key)
cache.set(key, value, ttl)
return value
Write-Through: Write to cache and DB simultaneously. Cache always consistent with DB. Higher write latency.
Write-Behind (Write-Back): Write to cache, asynchronously flush to DB. Lower write latency but risk of data loss on cache failure.
Read-Through: Cache sits in front of DB. Application always reads from cache; cache fetches from DB on miss. Transparent to application.
Cache Eviction Policies:
- LRU (Least Recently Used) — evict the least recently accessed
- LFU (Least Frequently Used) — evict the least accessed overall
- TTL — expire after a set time
- FIFO — evict oldest entry
Cache Stampede (Thundering Herd): Cache entry expires, N concurrent requests all miss and hit the DB simultaneously. Fix: probabilistic early expiration, mutex/distributed lock on miss, short stale-while-revalidate window.
CDN (Content Delivery Network)
A geographically distributed network of proxy servers that cache static and dynamic content close to users. Reduces latency by serving from the nearest edge node instead of origin.
Push CDN: You push content to CDN upfront. Good for large static assets that don't change often.
Pull CDN: CDN fetches from origin on first request, caches, serves subsequent requests from cache. Good for dynamic content.
Use CDNs for: images, JS/CSS, fonts, video streaming, and increasingly for API responses with short TTLs.
Database
ACID properties:
- Atomicity — transaction is all-or-nothing
- Consistency — transaction brings DB from one valid state to another
- Isolation — concurrent transactions appear serialized
- Durability — committed data persists even on crash
BASE properties (NoSQL):
- Basically Available — system always responds (maybe stale)
- Soft State — state may change over time without input
- Eventually Consistent — converges to consistent state eventually
Relational DB use cases: Strong consistency, complex queries, normalized data, ACID transactions.
NoSQL use cases:
- Key-Value (Redis, DynamoDB): Simple lookups, sessions, caching
- Document (MongoDB): Semi-structured data, flexible schema
- Wide-Column (Cassandra, HBase): Write-heavy, time series, IoT
- Graph (Neo4j): Social networks, recommendation engines
Indexes
An index is a data structure that speeds up reads at the cost of additional write overhead and storage.
B-Tree index (most relational DBs): Balanced tree, O(log N) lookup. Good for range queries and exact lookups. Default index type.
Hash index: O(1) exact lookup. Cannot do range queries.
Composite index: Index on multiple columns. The order matters — (user_id, created_at) efficiently supports WHERE user_id = ? ORDER BY created_at but not WHERE created_at = ? alone.
Covering index: Contains all columns needed by a query — DB can satisfy the query from the index without accessing the main table.
Index selectivity: A high-selectivity index (many unique values) is more efficient. Indexing a boolean column is usually useless.
Message Queues
Decouple producers from consumers. Producer sends a message; consumer processes it asynchronously.
Point-to-Point: One producer, one consumer per message. RabbitMQ queues.
Publish-Subscribe: One producer, many consumers all receive the message. Kafka topics.
Use cases: Async processing (sending emails, image resizing), traffic spike buffering, event-driven architectures, reliable inter-service communication.
At-least-once delivery: Message may be delivered multiple times on retry. Consumer must be idempotent.
At-most-once: May be lost, never duplicated. Good for metrics/analytics.
Exactly-once: No loss, no duplicates. Hardest to achieve; requires coordination.
Rate Limiting
Prevent abuse and protect services from being overwhelmed.
Fixed Window: Count requests in fixed time windows (e.g., 100 req/min). Burst problem: 100 at 0:59, 100 at 1:01 = 200 in 2 seconds.
Sliding Window: Count requests in a rolling window. More accurate but higher memory usage.
Token Bucket: Bucket fills at fixed rate; each request consumes one token. Allows bursts up to bucket size. Used in most APIs (AWS, Stripe).
Leaky Bucket: Requests enter a queue (bucket) that drains at a fixed rate. Smooths bursts but can drop excess.
Implementation: Redis with per-user counters and TTL is the standard approach.
Consistent Hashing
Used in distributed caches and databases to minimize data movement when nodes join/leave.
Arrange nodes on a virtual ring. Each node is responsible for keys between itself and its predecessor. Add a node: only neighboring keys rebalance. Remove a node: only its keys rebalance. Without consistent hashing, adding a node remaps almost all keys.
Virtual nodes: Each physical node has multiple positions on the ring for better load distribution.
Used by: DynamoDB, Cassandra, Memcached clusters.
API Design
REST: Resource-oriented URLs, stateless, standard HTTP verbs. Simple, widely supported, cacheable.
GraphQL: Client specifies exactly what data it needs. Single endpoint. Reduces over-fetching. Good for complex client requirements. More complex server implementation.
gRPC: Binary protocol (Protobuf), HTTP/2. Strongly typed, bi-directional streaming, faster than JSON. Good for internal microservices.
Idempotency: An operation that produces the same result no matter how many times it's executed. GET, PUT, DELETE are idempotent. POST is not by default.
Pagination: Never return unbounded result sets.
- Offset pagination:
?page=3&limit=20— simple but slow for large offsets - Cursor pagination:
?after=cursor_id&limit=20— O(1) regardless of position
Replication vs Sharding
Replication: Same data on multiple nodes. Purpose: fault tolerance and read scalability. All replicas have all data.
Sharding (Partitioning): Different data on different nodes. Purpose: write scalability when data exceeds single node capacity. Each shard has a subset of data.
Most large-scale systems use both: shard the data, then replicate each shard.
Key Takeaways
- Latency and throughput are related but distinct — optimize for the one that matters to your users
- Horizontal scaling requires stateless services; externalize state to caches and databases
- Caching is the highest-leverage optimization — pick the right eviction policy and invalidation strategy
- Load balancing algorithms matter — least connections beats round robin for variable workloads
- Consistent hashing minimizes data movement in distributed stores
- Message queues decouple write and processing — essential for handling traffic spikes
- CAP and BASE aren't just theory — they dictate which database you choose for a given access pattern
System design isn't about memorizing blueprints. It's about understanding these fundamentals deeply enough to reason about trade-offs and make deliberate decisions for your specific constraints.