Distributed systems fundamentals
Once two processes talk over a network, such as a service and its database, you have a distributed system. A few rules set it apart from single-machine programming, and they explain many patterns in earlier lessons.
Partial failure
On one machine, a program usually either works or crashes. In a distributed system, some parts fail while others keep running. This is partial failure.
When you send a request and get no response, you cannot tell which of these happened:
- the request was lost before it arrived;
- the remote side is down or very slow;
- the work was done but the response was lost.
They look the same to the caller. That is why you need timeouts, and why retries must be idempotent: after a timeout you don’t know whether the operation happened.
The network is not reliable
Messages can be delayed, dropped, duplicated or reordered. A network partition can split nodes into groups that can’t reach each other, while each group keeps running and believes the others have failed.
Clocks and ordering
Each machine has its own clock, and clocks drift. NTP corrects them imperfectly, and a correction can make a clock jump backwards. So wall-clock timestamps from different machines can’t reliably tell you which event came first. “Last write wins” by timestamp can silently discard the write that really came later.
When order matters, use something that doesn’t depend on clocks:
- a single authority that assigns order, such as a database sequence or one partition of a log;
- version numbers that increase with every change, as in optimistic locking;
- logical clocks, which order events by causality rather than time.
For measuring durations on one machine, use a monotonic clock, which never goes backwards.
Consistency models
- Strong consistency (linearisability): once a write completes, every later read sees it. The system behaves like a single copy of the data.
- Eventual consistency: if writes stop, all replicas converge to the same value, but in the meantime a read may return stale data.
Strong consistency needs coordination between nodes, which costs latency and availability. Eventual consistency is cheaper and survives failures but pushes complexity into the application. Lagging read replicas are a form of it.
CAP in practical terms
The CAP theorem says that during a network partition, a system must choose between consistency (refuse or delay requests it can’t answer correctly) and availability (answer anyway, possibly with stale or conflicting data). Partitions are not optional, so the real question is what each operation should do when one happens.
Partition between regions:
CP choice: reject writes in the minority side
AP choice: accept writes on both sides, reconcile later
The choice can differ per feature: balances usually want consistency; a like counter can favour availability.
Consensus
Some problems need nodes to agree on a value despite failures: who is leader, who holds a lock, what the configuration is. Algorithms such as Raft and Paxos solve this, but they are easy to get subtly wrong, with bugs that appear only under rare failure timings. Use a proven system such as etcd or ZooKeeper, or a database that provides the guarantee, rather than building your own.
Sagas instead of distributed transactions
A transaction spanning several services, such as placing an order, charging a card and reserving stock, could use two-phase commit. But it blocks when the coordinator fails and couples services tightly, so it is rarely used between services.
A saga is the common alternative: a sequence of local transactions, each in one service, with a compensating action for each step that undoes it if a later step fails. If stock reservation fails, the saga refunds the card and cancels the order. Other requests can see intermediate states, and every step and compensation must be idempotent.
Habits
- Treat every remote call as able to fail, hang or succeed without telling you, and make it idempotent.
- Don’t order events across machines by wall-clock time.
- Decide per feature whether stale reads are acceptable.
- Use existing consensus systems for leader election and locks.
- Prefer sagas with compensating actions over cross-service transactions.