System Design

System Design Important Terminology

August 20, 2024

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.

VA
Vishal
Aggarwal

Full Stack Developer

Ask about Vishal ✦