Design scalable, reliable, and fault-tolerant distributed systems using proven patterns and consistency models.
Purpose
Distributed systems are the foundation of modern cloud-native applications. Understanding fundamental trade-offs (CAP theorem, PACELC), consistency models, replication patterns, and resilience strategies is essential for building systems that scale globally while maintaining correctness and availability.
When to Use This Skill
Apply when:
Designing microservices architectures with multiple services
Building systems that must scale across multiple datacenters or regions
Choosing between consistency vs availability during network partitions
Designing partition-tolerant systems with proper consistency guarantees
Building resilient services with circuit breakers, bulkheads, retries
Implementing service discovery and inter-service communication
Core Concepts
CAP Theorem Fundamentals
CAP Theorem: In a distributed system experiencing a network partition, choose between Consistency (C) or Availability (A). Partition tolerance (P) is mandatory.
Network partitions WILL occur → Always design for P
During partition:
├─ CP (Consistency + Partition Tolerance)
│ Use when: Financial transactions, inventory, seat booking
│ Trade-off: System unavailable during partition
│ Examples: HBase, MongoDB (default), etcd
│
└─ AP (Availability + Partition Tolerance)
Use when: Social media, caching, analytics, shopping carts
Trade-off: Stale reads possible, conflicts need resolution
Examples: Cassandra, DynamoDB, Riak
PACELC: Extends CAP to consider normal operations (no partition).
If Partition: Choose Availability (A) or Consistency (C)
Else (normal): Choose Latency (L) or Consistency (C)
├─ Single region writes? → Leader-Follower
├─ Multi-region writes + conflicts OK? → Multi-Leader
├─ Multi-region writes + no conflicts? → Leader-Follower with failover
└─ Maximum availability? → Leaderless (quorum)
Choosing Partitioning Strategy
├─ Need range scans? → Range Partitioning (risk: hot spots)
├─ Data residency requirements? → Geographic Partitioning
└─ Default? → Hash Partitioning (consistent hashing)
Quick Reference Tables
CAP/PACELC System Comparison
| System | If Partition | Else (Normal) | Use Case |
For Kubernetes deployment: See kubernetes-operations skill for pod anti-affinity, service mesh
For infrastructure: See infrastructure-as-code skill for deploying distributed systems
For databases: See databases-sql and databases-nosql for replication configuration
For messaging: See message-queues skill for event-driven architectures, saga orchestration
For monitoring: See observability skill for distributed tracing, monitoring patterns
For testing: See performance-engineering skill for load testing distributed systems
For security: See security-hardening skill for mTLS, service authentication
Common Patterns
Multi-Datacenter Pattern
1. Choose replication: Multi-leader or Leaderless
2. Partition data geographically
3. Implement conflict resolution (LWW, vector clocks, app-specific)
4. Monitor replication lag
5. Add circuit breakers between datacenters
Event-Driven Saga Pattern
1. Define saga steps and compensating actions
2. Choose choreography (events) or orchestration (coordinator)
3. Implement idempotent handlers (retries safe)
4. Publish events with outbox pattern (transactional)
5. Monitor saga progress and timeouts
High-Availability Pattern
1. Use leaderless replication (N=5, W=3, R=2)
2. Partition with consistent hashing
3. Add circuit breakers for failing nodes
4. Implement read repair and anti-entropy
5. Monitor quorum health
Best Practices
Design for Failure:
Network partitions will occur - always design for partition tolerance
Use timeouts, retries with exponential backoff
Implement circuit breakers to prevent cascading failures
Test chaos engineering scenarios (partition nodes, inject latency)
Choose Consistency Carefully:
Default to eventual consistency, strengthen only where needed
Strong consistency has real costs (latency, availability)
Use bounded staleness for middle ground
Idempotency is Critical:
Design operations to be safely retryable
Use unique request IDs for deduplication
Essential for saga compensating transactions
Monitor and Observe:
Distributed tracing with correlation IDs
Monitor replication lag, quorum health
Alert on circuit breaker state changes
Track saga progress and failures
Partition Strategically:
Hash partitioning for even distribution
Range partitioning for range queries (monitor hot spots)
Geographic partitioning for compliance, latency
Version Everything:
Event schemas evolve - use versioning
API versioning for service compatibility
Database schema migrations in distributed systems
Anti-Patterns to Avoid
Distributed Monolith:
Microservices with tight coupling
Shared database across services
Fix: Database per service, async communication
Two-Phase Commit (2PC) Overuse:
Slow, blocking, reduces availability
Fix: Use saga pattern for distributed transactions