Distributed Systems
Designing at scale means trading off consistency, availability, and partition tolerance, and using replication, sharding, and consensus to meet goals.
๐ Overview
Distributed systems coordinate multiple nodes over unreliable networks. Latency, partial failures, and partitions are the norm, requiring careful consistency and availability designs.
๐ Fundamentals
- CAP: choose consistency vs availability under partitions
- Consensus: Raft/Paxos to agree under failures
- Replication & Sharding: scale reads and writes
๐ ๏ธ Guides
Consistency Models
- Strong: linearizable reads/writes
- Eventual: replicas converge without strict ordering
- Causal: respect happens-before relationships
Design Patterns
- Leader-based replication with log
- Quorums (R/W) for availability/consistency trade-offs
- Idempotency and retries for at-least-once delivery
๐ Applications
- Datastores: key-value stores, document DBs
- Queues/streams: distributed logs and messaging
- Coordination: service discovery, leader election
๐งช Examples
- Raft log replication and state machine apply
- Consistent hashing for scalable partitioning
- Two-phase commit vs saga for distributed transactions
โ Frequently Asked Questions
1) Does CAP forbid strong systems?
No; under no partition, you can have both. Under partition, you must choose.
No; under no partition, you can have both. Under partition, you must choose.
2) Paxos vs Raft?
Equivalent properties; Raft is often simpler to implement.
Equivalent properties; Raft is often simpler to implement.
3) Exactly-once delivery?
Achieved via idempotency and deduplication semantics.
Achieved via idempotency and deduplication semantics.
4) Clock sync?
Use logical/Hybrid clocks; do not depend on perfect NTP.
Use logical/Hybrid clocks; do not depend on perfect NTP.
5) Hot partitions?
Mitigate with better keys, load-aware routing, or rebalancing.
Mitigate with better keys, load-aware routing, or rebalancing.
6) Multi-region writes?
Conflict-free replicated data types (CRDTs) or operational transforms.
Conflict-free replicated data types (CRDTs) or operational transforms.
7) Testing failures?
Chaos engineering to exercise partitions and node loss.
Chaos engineering to exercise partitions and node loss.
8) Backpressure?
Control producer rates to protect downstream services.
Control producer rates to protect downstream services.
9) Data migrations?
Dual writing and read routing during cutovers.
Dual writing and read routing during cutovers.
10) Observability?
Tracing, metrics, and logs across services and requests.
Tracing, metrics, and logs across services and requests.