Building distributed systems is fundamentally about managing failure. A single machine either works or it doesn't. In a distributed system, parts of it are always failing — the challenge is making the whole system behave correctly despite that.
This article covers the foundational concepts I had to internalize before I could reason clearly about distributed systems.
Why Distribution?
A single machine has hard limits: CPU, memory, disk throughput, and network bandwidth. Beyond a certain scale, you must distribute work across multiple machines. But distribution introduces problems that don't exist on a single node:
- Partial failures — some nodes are up, some are down
- Network unreliability — packets get lost, delayed, or duplicated
- No shared clock — you can't assume events on different nodes happened in a consistent order
- No shared memory — nodes must communicate explicitly
The CAP Theorem
CAP states that a distributed system can guarantee at most two of:
- Consistency (C) — every read receives the most recent write or an error
- Availability (A) — every request receives a response (not necessarily the latest data)
- Partition Tolerance (P) — the system continues operating despite network partitions
Since network partitions are a reality you cannot avoid, you're always choosing between CP or AP behavior when a partition occurs.
CP systems (e.g., HBase, Zookeeper) — will refuse requests rather than return stale data. Good for financial transactions, configuration stores.
AP systems (e.g., Cassandra, DynamoDB) — will serve potentially stale data rather than return errors. Good for shopping carts, social feeds, DNS.
The CAP theorem is often misunderstood. In normal operation (no partition), you can have both consistency and availability. The trade-off only applies under partition.
Consistency Models
CAP's "consistency" is linearizability — the strongest model. In practice, systems offer a spectrum:
Linearizability (Strong Consistency) — every operation appears instantaneous and in a globally consistent order. Reads always return the latest write. Most expensive.
Sequential Consistency — all nodes see operations in the same order, but not necessarily real-time. Weaker than linearizability.
Causal Consistency — if A causally depends on B, all nodes see B before A. Unrelated operations can be seen in any order.
Eventual Consistency — given no new writes, all replicas will converge to the same value. Eventual doesn't mean slow — it means the guarantee is eventual, not immediate.
Read-your-writes — a client always reads its own writes. Weaker than strong consistency but often sufficient for user-facing features.
Replication
Replication keeps multiple copies of data on different nodes for fault tolerance and read throughput.
Single-leader replication — one node accepts writes (leader), replicates to followers. Followers serve reads. Simple but leader is a bottleneck and SPOF.
Multi-leader replication — multiple leaders accept writes. Used in geo-distributed systems. Requires conflict resolution.
Leaderless replication — any node accepts writes and reads (e.g., Dynamo-style). Uses quorums:
- Write to W nodes, read from R nodes
- If W + R > N (total replicas), reads will overlap with writes — strong consistency
- Lower W/R = higher availability, lower consistency
Consensus
Consensus is the problem of getting a group of nodes to agree on a value. It's needed for leader election, distributed transactions, and replicated state machines.
Paxos — the original consensus algorithm. Correct but notoriously hard to understand and implement.
Raft — designed to be understandable. Separates leader election, log replication, and safety. Used in etcd, CockroachDB.
Core guarantee: once a value is decided, all correct nodes will eventually agree on it even if some nodes crash.
Distributed Transactions
How do you update two databases atomically — either both succeed or both fail?
Two-Phase Commit (2PC):
- Coordinator asks all participants to "prepare" (lock resources, check feasibility)
- If all say "yes", coordinator sends "commit"; if any says "no", sends "abort"
Problem: if the coordinator crashes after "prepare" but before "commit", participants are stuck holding locks forever. 2PC is a blocking protocol.
Saga Pattern — split the transaction into a sequence of local transactions. Each step publishes an event. If a step fails, compensating transactions roll back prior steps. No distributed locking required, but you must design compensating actions.
Ordering and Clocks
Without a global clock, how do you know which event happened first?
Lamport Timestamps — each node maintains a logical clock counter. On any event, increment the counter. On message send, include the counter. On message receive, take max(local, received) + 1. Lamport clocks establish a partial order: if A → B then L(A) < L(B), but the converse isn't guaranteed.
Vector Clocks — each node maintains a vector of counters (one per node). Captures causality precisely — you can determine if two events are concurrent or causally related.
Hybrid Logical Clocks (HLC) — combine physical time with logical clocks. Used in CockroachDB for reading at a consistent snapshot without locking.
Failure Detection
In a distributed system, you can't reliably distinguish between a crashed node and a slow one. Failure detectors are inherently probabilistic.
Heartbeats — nodes send periodic pings. If no heartbeat in timeout T, declare dead. Simple but tuning T is tricky: too small → false positives, too large → slow detection.
Phi Accrual Detector — instead of a binary alive/dead decision, outputs a suspicion level φ. Higher φ = more confident the node is down. Used in Cassandra and Akka.
Gossip Protocols — nodes periodically exchange state with random peers. Information spreads exponentially. Highly resilient to node failures. Used for membership and failure detection in Cassandra.
Partitioning (Sharding)
Partitioning distributes data across nodes so no single node holds all data.
Range partitioning — split by key range. Easy to range-scan but can create hot spots if access patterns are skewed.
Hash partitioning — hash the key, assign to a bucket. Distributes load evenly but makes range queries expensive.
Consistent hashing — arrange nodes on a virtual ring. Keys map to the nearest node clockwise. Adding/removing nodes moves only O(1/N) keys. Used in DynamoDB, Cassandra.
Virtual nodes (vnodes) — each physical node claims multiple positions on the ring. Better load distribution and simpler rebalancing when nodes join/leave.
CRDT: Conflict-Free Replicated Data Types
CRDTs are data structures designed for eventual consistency — any two replicas can be merged and the result is always the same regardless of merge order.
Examples: G-Counter (grow-only), PN-Counter (increment and decrement), LWW-Register (last-write-wins), OR-Set (observed-remove set).
CRDTs are used in collaborative editing (Google Docs), shopping carts, and distributed caches where conflicts are expected and you want automatic resolution.
Key Takeaways
- Distributed systems fail partially — design for it, not around it
- CAP: you're choosing between CP and AP under partition, not globally
- Consensus (Raft/Paxos) is needed for anything requiring agreement across nodes
- 2PC is blocking; Saga is not — choose based on your failure tolerance
- Logical clocks (Lamport, Vector) track causality without a global clock
- Consistent hashing minimizes rebalancing overhead when topology changes
The deeper you go into distributed systems, the more you appreciate that there are no free lunches — every design choice trades off latency, throughput, consistency, or availability. The job of a distributed systems engineer is to make those trade-offs deliberately.