Skip to content

Side 58

Distributed
Systems

A study of computation spread across unreliable machines. Distributed systems replace shared memory with messages, introducing delay, partial failure, uncertain ordering and the need for explicit coordination.

node→message→failure→coordination→consistency
06system lenses
05failure modes
05coordination problems
58Side

Distribution changes the failure model.

One process can no longer directly know the state of another; it only receives messages that may be delayed, duplicated or lost.

01 · Node

Independent computer or process.

Own memory + own clock.

Nodes do not share instantaneous state.

02 · Message

Communication crosses a network.

Delay is variable.

A missing reply does not reveal whether the peer failed or the network is slow.

03 · State

Copies can diverge.

Who knows the latest value?

Replication creates availability and coordination problems simultaneously.

04 · Failure

Some parts fail while others continue.

Partial failure.

Distributed systems must operate without assuming all components share fate.

05 · Recovery

Rejoin without corrupting state.

Replay, reconcile, elect?

Recovery logic is part of the protocol, not an afterthought.

There is no perfectly shared clock.

Distributed systems often need to reason about order without relying on synchronized wall time.

Clock drift

Physical clocks diverge.

Synchronization reduces but does not eliminate uncertainty.

Happens-before

Causal order without global time.

If one event can influence another, the system can order them causally.

Logical clock

Represent ordering abstractly.

Lamport clocks track a partial temporal relation among events.

Vector clock

Track causal histories.

Vector clocks can distinguish concurrent updates from causally ordered ones.

Timeout

Bound waiting.

Timeouts turn uncertainty into a decision point without proving failure.

Lease

Grant authority for bounded time.

Leases depend on carefully managed clock assumptions.

Replication trades coordination for availability and performance.

Copies increase resilience and read capacity, but writes must be propagated and reconciled.

Leader

One node orders writes.

Followers replicate the leader’s log or state.

Multi-leader

Several nodes accept writes.

Useful across regions but creates conflict-resolution work.

Leaderless

Clients write to several replicas.

Quorums and repair mechanisms reconcile divergent copies.

Lag

Replicas are temporarily behind.

Read behavior depends on whether stale results are acceptable.

Conflict

Concurrent writes disagree.

Systems need deterministic merge or application-level resolution rules.

Consensus creates one agreed history despite failures.

Protocols such as Raft and Paxos coordinate nodes on an ordered sequence of decisions.

Election

Choose a coordinator.

Nodes must avoid accepting conflicting leaders for the same logical term.

Quorum

Require overlapping majorities.

Overlap prevents two independent decisions from both appearing committed.

Log

Replicate an ordered history.

State machines can replay the same committed commands deterministically.

Commit

Decide when an entry is durable enough.

Consensus separates proposed work from work known to be committed.

Safety

Never decide conflicting values.

Safety may be preserved even when progress temporarily stops.

Liveness

Eventually make progress.

Progress depends on enough functioning nodes and communication.

Consistency defines what clients are allowed to observe.

Different guarantees make different promises about ordering and freshness.

ModelGuaranteeCost
LinearizableOperations appear to occur atomically in real-time orderCoordination and latency
SequentialAll clients see one order consistent with program orderStill coordinated, weaker than real-time
Read-your-writesA client sees its own completed updatesSession routing/state
EventualReplicas converge if updates stopTemporary divergence allowed
CausalCausally related operations preserve orderMetadata and dependency tracking

Reliability comes from designing around failure, not assuming it away.

Distributed systems need idempotence, retries, redundancy and clear ownership of uncertain work.

Retry

Repeat transiently failed operations, but make repeated execution safe where possible.

Idempotence

Ensure duplicate requests do not create duplicate effects.

Backoff

Slow retries to avoid amplifying an overloaded dependency.

Circuit breaker

Stop sending work to a dependency that is repeatedly failing.

Dead letter

Preserve work that cannot be processed automatically for later inspection.

Observability

Correlate logs, metrics and traces across nodes to reconstruct failure paths.

Designing Data-Intensive ApplicationsKleppmann · distributed data systems
Distributed SystemsTanenbaum & van Steen · broad foundation
In Search of an Understandable Consensus AlgorithmOngaro & Ousterhout · Raft
Distributed AlgorithmsNancy Lynch · formal foundations