LearnThatStack Ace your next interview

Distributed Systems.
Interview cheat sheet.

Quick reference for Distributed Systems - sectioned for fast scanning. Skim the part you're shaky on, walk in confident.

System Design Concepts 17-section reference ~6 min read

Summary

Comprehensive guide to distributed systems for technical interviews. Covers fundamental concepts, CAP theorem, consensus algorithms, scaling strategies, fault tolerance, and real-world system designs. Essential knowledge for architecting systems that span multiple machines and handle millions of users.

1. Core Concepts & Definitions

What is a Distributed System?

  • Definition: Collection of independent computers that appears as a single coherent system
  • Key Properties:
    • No shared memory
    • Communication via network
    • Independent failures
    • No global clock

Why Distributed Systems?

  • Scalability: Handle more load by adding machines
  • Reliability: Continue operating despite failures
  • Performance: Geographic distribution reduces latency
  • Cost: Commodity hardware vs expensive supercomputers

2. Fundamental Principles

CAP Theorem

Choose 2 out of 3:

  • Consistency (C): All nodes see the same data simultaneously
  • Availability (A): System remains operational
  • Partition Tolerance (P): System continues despite network failures

Common Choices:

  • CP Systems: MongoDB, HBase, Redis
  • AP Systems: Cassandra, DynamoDB, CouchDB
  • CA Systems: Traditional RDBMS (single node)

PACELC Theorem

Extension of CAP: "If Partition, then tradeoff between Availability and Consistency; Else, tradeoff between Latency and Consistency"

Consistency Models

  1. Strong Consistency: All nodes see same data instantly
  2. Eventual Consistency: Nodes converge to same state eventually
  3. Weak Consistency: No guarantees about when data syncs
  4. Read-after-Write: User sees their own writes immediately
  5. Monotonic Read: Once seen, data doesn't go backward

3. System Design Patterns

Load Balancing

Types:

  • Round Robin: Requests distributed sequentially
  • Least Connections: Route to server with fewest active connections
  • IP Hash: Route based on client IP
  • Geographic: Route to nearest datacenter

Implementation:

# Simple round-robin example
class LoadBalancer:
    def __init__(self, servers):
        self.servers = servers
        self.current = 0
    
    def get_server(self):
        server = self.servers[self.current]
        self.current = (self.current + 1) % len(self.servers)
        return server

Caching Strategies

  1. Cache-Aside (Lazy Loading)

    • Read: Check cache → miss → load from DB → update cache
    • Write: Write to DB → invalidate cache
  2. Write-Through

    • Write to cache and DB simultaneously
    • Ensures consistency but higher latency
  3. Write-Behind (Write-Back)

    • Write to cache → async write to DB
    • Better performance but risk of data loss

Message Queues

Use Cases:

  • Decoupling components
  • Handling traffic spikes
  • Ensuring reliability

Popular Systems:

  • RabbitMQ: Feature-rich, AMQP protocol
  • Kafka: High-throughput, distributed log
  • Redis Pub/Sub: Simple, in-memory
  • SQS: AWS managed service

4. Data Storage & Management

Database Replication

Master-Slave (Primary-Replica)

  • Master handles writes
  • Slaves handle reads
  • Provides read scalability

Master-Master

  • Multiple masters accept writes
  • Requires conflict resolution
  • Higher availability

Database Sharding

Strategies:

  1. Range-based: Shard by data range (A-M, N-Z)
  2. Hash-based: Shard by hash(key) % N
  3. Geographic: Shard by location
  4. Directory-based: Lookup service for shard location

Example:

# Simple hash-based sharding
def get_shard(key, num_shards):
    return hash(key) % num_shards

Consistent Hashing

  • Minimizes data movement when nodes added/removed
  • Used in: Cassandra, DynamoDB, Memcached
# Simplified consistent hashing
class ConsistentHash:
    def __init__(self, nodes, virtual_nodes=150):
        self.ring = {}
        for node in nodes:
            for i in range(virtual_nodes):
                key = hash(f"{node}:{i}")
                self.ring[key] = node
        self.sorted_keys = sorted(self.ring.keys())
    
    def get_node(self, key):
        if not self.ring:
            return None
        hash_key = hash(key)
        # Find first node clockwise
        for node_key in self.sorted_keys:
            if hash_key <= node_key:
                return self.ring[node_key]
        return self.ring[self.sorted_keys[0]]

5. Distributed Algorithms

Consensus Algorithms

Paxos

  • Ensures agreement among distributed nodes
  • Complex but proven correct
  • Used in: Google's Chubby

Raft

  • Easier to understand alternative to Paxos
  • Leader election + log replication
  • Used in: etcd, Consul

Key Concepts:

  • Leader Election: One node coordinates
  • Log Replication: Ensures consistency
  • Term Numbers: Logical time periods

Clock Synchronization

Lamport Timestamps

  • Logical ordering of events
  • Each process maintains counter
  • Increment on local event, update on message receive

Vector Clocks

  • Tracks causality between events
  • Each node has vector of counters
  • Detects concurrent events

6. Microservices Architecture

Service Discovery

Approaches:

  1. Client-side: Client queries registry (Eureka)
  2. Server-side: Load balancer queries registry
  3. Service Mesh: Sidecar proxy handles discovery (Istio)

API Gateway

Functions:

  • Request routing
  • Authentication/authorization
  • Rate limiting
  • Request/response transformation
  • Monitoring

Circuit Breaker Pattern

class CircuitBreaker:
    def __init__(self, failure_threshold=5, timeout=60):
        self.failure_threshold = failure_threshold
        self.timeout = timeout
        self.failure_count = 0
        self.last_failure_time = None
        self.state = "CLOSED"  # CLOSED, OPEN, HALF_OPEN

7. Fault Tolerance & Reliability

Failure Types

  1. Crash Failure: Node stops responding
  2. Omission Failure: Messages lost
  3. Timing Failure: Response too slow
  4. Byzantine Failure: Arbitrary/malicious behavior

Handling Failures

  • Timeouts: Detect unresponsive services
  • Retries: Handle transient failures
  • Fallbacks: Graceful degradation
  • Bulkheads: Isolate failures
  • Health Checks: Proactive monitoring

Replication Strategies

  1. Single-leader: One master, multiple replicas
  2. Multi-leader: Multiple masters, conflict resolution needed
  3. Leaderless: All nodes equal (Dynamo-style)

8. Scaling Strategies

Horizontal vs Vertical Scaling

Vertical (Scale-up)

  • Add more power (CPU, RAM) to existing machine
  • Limited by hardware constraints
  • No code changes required

Horizontal (Scale-out)

  • Add more machines
  • Better fault tolerance
  • Requires distributed system design

Database Scaling

  1. Read Replicas: Distribute read load
  2. Sharding: Distribute data across nodes
  3. Federation: Split databases by function
  4. Denormalization: Trade space for query performance

9. Communication Protocols

REST vs RPC

REST

  • HTTP-based, stateless
  • Resource-oriented
  • Cacheable responses

RPC (gRPC, Thrift)

  • Binary protocols
  • Strongly typed
  • Better performance
  • Bidirectional streaming

Message Patterns

  1. Request-Reply: Synchronous communication
  2. Publish-Subscribe: One-to-many broadcast
  3. Message Queue: Asynchronous, decoupled

10. Real-World System Examples

URL Shortener

Components:

  • API Gateway
  • Application servers
  • Cache layer (Redis)
  • Database (sharded)
  • Analytics service

Key Decisions:

  • Base62 encoding for short URLs
  • Custom URL support
  • Cache popular URLs
  • Geographic distribution

Social Media Feed

Challenges:

  • Real-time updates
  • Celebrity problem (millions of followers)
  • Personalization

Solutions:

  • Push model for normal users
  • Pull model for celebrities
  • Hybrid approach
  • Pre-computed timelines

Distributed File Storage

Components:

  • Metadata servers
  • Chunk servers
  • Client library

Key Features:

  • File chunking (64MB blocks)
  • Replication (3x typically)
  • Checksums for integrity
  • Geographic distribution

11. Performance Optimization

Caching Levels

  1. Browser Cache: Client-side
  2. CDN: Geographic distribution
  3. Application Cache: In-memory (Redis)
  4. Database Cache: Query results

Database Optimization

  • Indexing: B-trees, hash indexes
  • Query Optimization: Explain plans
  • Connection Pooling: Reuse connections
  • Denormalization: Trade storage for speed

12. Monitoring & Observability

Key Metrics

  • Latency: Response time (p50, p95, p99)
  • Traffic: Requests per second
  • Errors: Error rate and types
  • Saturation: Resource utilization

Distributed Tracing

  • Track requests across services
  • Tools: Jaeger, Zipkin, AWS X-Ray
  • Correlation IDs for request tracking

13. Security Considerations

Authentication & Authorization

  • OAuth 2.0: Delegated authorization
  • JWT: Stateless authentication
  • API Keys: Service-to-service auth
  • mTLS: Mutual TLS for zero-trust

Data Security

  • Encryption at Rest: Protect stored data
  • Encryption in Transit: TLS/SSL
  • Key Management: Rotate regularly
  • Data Masking: Protect sensitive info

14. Interview Tips

System Design Approach

  1. Clarify Requirements

    • Functional requirements
    • Non-functional (scale, performance)
    • Constraints
  2. Estimate Scale

    • Users, requests/second
    • Storage requirements
    • Bandwidth needs
  3. Design High-Level

    • Major components
    • Data flow
    • API design
  4. Deep Dive

    • Data models
    • Algorithm choices
    • Optimization strategies
  5. Handle Failures

    • Single points of failure
    • Data loss scenarios
    • Monitoring/alerting

Common Pitfalls

  • Over-engineering for current scale
  • Ignoring data consistency
  • Not considering failures
  • Missing bottlenecks
  • Forgetting about monitoring

Key Concepts to Master

  • CAP theorem and trade-offs
  • Consistent hashing
  • Database sharding strategies
  • Caching patterns
  • Load balancing algorithms
  • Message queue use cases
  • Microservices patterns
  • Consensus algorithms basics
  • Scaling techniques
  • Common system designs

Quick Reference

Latency Numbers (2020s)

  • L1 cache: 0.5 ns
  • L2 cache: 7 ns
  • RAM access: 100 ns
  • SSD read: 150 μs
  • HDD read: 10 ms
  • Network round trip (same region): 0.5 ms
  • Network round trip (across regions): 150 ms

Common Port Numbers

  • HTTP: 80
  • HTTPS: 443
  • MySQL: 3306
  • PostgreSQL: 5432
  • MongoDB: 27017
  • Redis: 6379
  • Elasticsearch: 9200
  • Kafka: 9092

HTTP Status Codes

  • 200: OK
  • 201: Created
  • 400: Bad Request
  • 401: Unauthorized
  • 403: Forbidden
  • 404: Not Found
  • 500: Internal Server Error
  • 502: Bad Gateway
  • 503: Service Unavailable

Key Takeaways

  • No Perfect Solution: Every design involves trade-offs
  • Failures Are Inevitable: Design for failure from the start
  • Scale Incrementally: Don't over-engineer for imaginary scale
  • Data Consistency: Understand the cost of consistency guarantees
  • Monitoring Critical: Can't fix what you can't measure
  • Caching is Key: But cache invalidation is hard
  • Async When Possible: Decouple components for resilience
  • Think in Layers: Separation of concerns improves maintainability
Found this useful? Pass it on.
Pro · $10/mo

The sheet is free. Pro goes deeper.

Pro opens the full question library behind every sheet, every refresher and a monthly AI allowance. One subscription, all formats.

Full question library All refreshers Cancel anytime