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
- Strong Consistency: All nodes see same data instantly
- Eventual Consistency: Nodes converge to same state eventually
- Weak Consistency: No guarantees about when data syncs
- Read-after-Write: User sees their own writes immediately
- 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
Cache-Aside (Lazy Loading)
- Read: Check cache → miss → load from DB → update cache
- Write: Write to DB → invalidate cache
Write-Through
- Write to cache and DB simultaneously
- Ensures consistency but higher latency
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:
- Range-based: Shard by data range (A-M, N-Z)
- Hash-based: Shard by hash(key) % N
- Geographic: Shard by location
- 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:
- Client-side: Client queries registry (Eureka)
- Server-side: Load balancer queries registry
- 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
- Crash Failure: Node stops responding
- Omission Failure: Messages lost
- Timing Failure: Response too slow
- 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
- Single-leader: One master, multiple replicas
- Multi-leader: Multiple masters, conflict resolution needed
- 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
- Read Replicas: Distribute read load
- Sharding: Distribute data across nodes
- Federation: Split databases by function
- 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
- Request-Reply: Synchronous communication
- Publish-Subscribe: One-to-many broadcast
- 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
- Browser Cache: Client-side
- CDN: Geographic distribution
- Application Cache: In-memory (Redis)
- 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
Clarify Requirements
- Functional requirements
- Non-functional (scale, performance)
- Constraints
Estimate Scale
- Users, requests/second
- Storage requirements
- Bandwidth needs
Design High-Level
- Major components
- Data flow
- API design
Deep Dive
- Data models
- Algorithm choices
- Optimization strategies
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