Skip to content
teach

Learning: Distributed Systems

Be able to choose and defend a consistency model for a system you are designing, and to reason about a production incident caused by a partial failure instead of treating the network as reliable.

Latest lesson: 19. Verifying a Claim

Success looks like

  • Given a system design, choose a consistency model and defend the trade-off against the alternatives.
  • Given an incident caused by a partial failure (a network partition, a slow node mistaken for a dead one), name the mechanism responsible.
  • Explain what a consensus protocol buys you and why it costs what it costs, without needing to prove its correctness from first principles.

Constraints

  • Assumes professional experience building a networked service; no prior formal distributed-systems study required.
  • Practical and operational emphasis: the cost and the trade-offs matter more than a formal proof of a protocol's correctness.

Out of scope

  • Concurrency inside one process (goroutines, threads, Send/Sync): owned by the language workspaces. This workspace starts where the process boundary is crossed.

The arc

Twelve stages, partial failure to a diagnosed incident to the mechanisms production systems actually use to survive one, and how to verify any of it actually holds. A stage takes several lessons and the boundaries are soft; what makes a stage done is the capability, not the lesson count.

Stage Lessons Covers Done when
1. Partial failure 0001 The one problem every later topic in this workspace is a response to Can explain why a network can't be treated as reliable
2. Time and order 0002 to 0003 Clocks, ordering, why a timeout is the only failure signal available Can explain why a timeout can't distinguish a slow node from a dead one
3. Consistency models 0004 to 0006 CAP, linearizability, sequential and eventual consistency Can choose and defend a consistency model for a stated design
4. Consensus 0007 to 0009 Raft, leader election, what consensus buys and why it costs what it costs Can explain what a consensus protocol buys and costs without proving its correctness
5. Diagnosing incidents 0010 Applying the mechanisms above to a real production incident Given an incident caused by partial failure, can name the mechanism responsible
6. Replication and quorums 0011 Leader-based vs. leaderless replication, read/write quorums, the R + W > N condition Can explain how a quorum guarantees a read sees the latest write, and choose R and W for a stated workload
7. Partitioning and sharding 0012 Consistent hashing, virtual nodes, hot-key replication, request routing Can explain why consistent hashing bounds the cost of a resize, and what a bare ring still needs to handle failure and hot keys
8. Conflict resolution under eventual consistency 0013 Last-write-wins and what it loses, version vectors, CRDTs Can distinguish detecting a conflict from resolving one, and explain how a CRDT merges without loss
9. Transactions across services 0014 to 0015 Two-phase commit and its blocking failure mode, sagas and compensation, the transactional outbox, idempotent consumers Can explain why 2PC blocks, what a saga trades away to avoid it, and how the outbox pattern avoids the dual-write problem
10. The resilience toolkit 0016 to 0017 Timeout budgets, retries with backoff and jitter, retry storms, circuit breakers, load shedding, backpressure Can explain what to do after a timeout fires, and match a failing dependency, an overloaded server, or a filling queue to the right mechanism
11. Observability of partial failure 0018 Spans, traces, propagation across a network boundary, span kinds Can explain how a trace ties one request's spans together across services, and why propagation must happen explicitly on every hop
12. Verifying a claim 0019 Jepsen's opaque-box, model-checked testing, chaos engineering's production fault injection, deterministic simulation's reproducibility Can explain what trade-off each verification technique makes, and what a given test result does and doesn't prove

Lessons

Work through these in order.

# Lesson Teaches
0001 Partial Failure The one problem every later topic in this workspace is a response to
0002 Clocks and Ordering Why wall-clock timestamps from different machines can't be trusted to order events, and how a logical clock orders them without needing synchronized time
0003 Timeouts as Failure Detectors Why every practical failure detector is built on a timeout, and the formal vocabulary for the accuracy-versus-speed trade-off that follows from it
0004 The CAP Theorem, Precisely What CAP actually proves, why partition tolerance was never optional, and the specific misreadings that make "pick two" the wrong way to state it
0005 Linearizability What linearizability actually guarantees, why it's the strongest common consistency model, and what it costs to provide during a partition
0006 Sequential and Eventual Consistency Two models weaker than linearizability, what each still guarantees, and how to choose among all three for a stated design
0007 What Consensus Is For The replicated-state-machine problem consensus protocols solve, and why it's the mechanism underneath a linearizable system's coordination
0008 Raft, Leader Election and Log Replication How Raft elects a single leader and replicates a log through it, the two mechanisms that turn the replicated-state-machine problem into something concrete
0009 What Consensus Costs Why every write pays a round trip to a majority, why a minority partition loses availability rather than consistency, and how to weigh that cost against what consensus buys
0010 Diagnosing a Production Incident Applying partial failure, failure detection, consistency models, and consensus to name the mechanism behind a real incident, instead of reasoning about "the network" or "consistency" in the abstract
0011 Replication and Quorums A consensus protocol isn't the only way to replicate data, and the quorum condition behind its cheaper alternative is a single overlap guarantee, not a vague notion of majority agreement
0012 Partitioning and Sharding Consistent hashing exists because naive hash-mod-N partitioning remaps almost everything the moment a node joins or leaves, and even consistent hashing needs virtual nodes before it stops dumping a failed node's whole load onto one unlucky neighbor
0013 Conflict Resolution Under Eventual Consistency Detecting that two writes were concurrent and deciding what to do about it are two different problems, and last-write-wins solves neither, it just picks a survivor and throws the loser away
0014 Two-Phase Commit and Sagas Two-phase commit looks like consensus (a coordinator, a vote, a commit), but it solves a different problem and pays a worse price for it, blocking forever if the coordinator dies mid-vote, which is exactly why production systems reach for sagas instead
0015 The Transactional Outbox and Idempotent Consumers Updating a database and publishing an event can't both happen atomically without a distributed transaction, so the outbox pattern sidesteps that entirely by writing the event as an ordinary row in the same local transaction, then paying for it with a message that might be sent twice
0016 Timeouts, Retries, and Backoff Lesson 3 established that a timeout is the only failure signal available, but firing one and retrying immediately is exactly what turns a recovering service's bad day into a pile-on, which is the specific problem backoff and jitter exist to prevent
0017 Circuit Breakers, Load Shedding, and Backpressure A retry with backoff and jitter still assumes the failing call is worth attempting at all, and a circuit breaker, load shedding, and backpressure are the three answers for when it stops being worth it, aimed at a failing dependency, an overloaded server, and an overwhelmed queue respectively
0018 Distributed Tracing and Correlation Every mechanism this workspace has covered so far tells you what a system does under failure, but none of it tells you which hop in a specific request actually failed, which is the one thing a trace is built to answer
0019 Verifying a Claim Jepsen, chaos engineering, and deterministic simulation all try to find a real bug before a customer does, but they make opposite trades between realism and reproducibility to get there, and knowing which is which is what tells you what a given test actually proved

Reference

  • Glossary: canonical terms for this topic
  • Resources: trusted sources
  • Time and Order: what happens-before can and cannot decide, what a Lamport clock refuses to tell you, and the failure-detector vocabulary a timeout is one instance of
  • Consistency Models: CAP stated as it was proved, the models side by side with what each guarantees, and which of them can answer at all during a partition
  • Consensus: the replicated-state-machine problem, Raft's five safety properties and the rules that produce them, and what a majority costs on writes and on reads
  • Diagnosing Incidents: four questions in order, the named anomaly vocabulary for describing a symptom precisely, and a symptom-to-mechanism table indexed to the sheets above

How this works

Each lesson is short and self-contained. Answer keys are collapsed: recall first, then open them. The real-world reps matter more than the reading, and spacing them out is the point. Anything still unclear at the end of a lesson is worth chasing to its primary source before moving on.

Table of contents