Distributed Systems Resources
Knowledge
- Article: "Fallacies of distributed computing", Wikipedia
The canonical list (originated by Peter Deutsch, extended by James Gosling) of assumptions that hold on one machine and quietly stop holding once a network sits between two of them. Use for: the vocabulary for what a network can do to you, before reasoning about any specific failure. - Paper: "Time, Clocks, and the Ordering of Events in a Distributed System", Leslie Lamport, 1978
The original paper defining the happens-before relation and logical clocks: a way to order events across machines without relying on synchronized physical clocks. Use for: why wall-clock timestamps from different machines can't be trusted to order events, and what to use instead. - Paper: "Virtual Time and Global States of Distributed Systems", Friedemann Mattern, 1988
One of the standard references for vector clocks, the mechanism that supplies the direction Lamport's clock condition deliberately does not: with a vector,ahappens-beforebif and only if the vectors compare, so concurrency becomes detectable. Use for: deciding whether a system needs to recognise concurrent events, and what that costs in message size. - Paper: "Impossibility of Distributed Consensus with One Faulty Process", Fischer, Lynch and Paterson, 1985
The FLP result: no deterministic algorithm in an asynchronous system can guarantee consensus with even one faulty process. Use for: why failure detectors exist at all, since they are the standard way of adding just enough assumption to escape this result without pretending the network is synchronous. - Paper: "Unreliable Failure Detectors for Reliable Distributed Systems", Chandra and Toueg, 1996
The paper formalizing failure detectors by their completeness and accuracy properties, and showing consensus is solvable even with a failure detector that makes infinitely many mistakes. Use for: the formal vocabulary behind why a practical failure detector (a timeout) trades accuracy for speed, and the bridge into what consensus protocols (stage 4) actually need from failure detection. - Site: "Consistency Models", Jepsen
An interactive, precisely-defined map of consistency models (linearizability, serializability, causal consistency, and more), with the guarantees and violations that distinguish each. Use for: pinning down exactly which consistency model a system is buying, rather than reasoning about "consistency" as one vague thing. - Paper: "Brewer's Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Services", Gilbert and Lynch, 2002
The proof, as opposed to the conjecture and the retrospective. Its definitions are narrower than the words: consistency means atomic, meaning linearizable, and availability means every request received by a non-failing node must terminate in a response, with no bound on how long. Use for: settling what CAP actually claims, especially when an argument turns on what "available" was supposed to mean. - Paper: "Linearizability: A Correctness Condition for Concurrent Objects", Herlihy and Wing, 1990
The paper that introduced linearizability, and with it strict serializability. Use for: the precise guarantee, and the fact that it is a single-object condition, so a system calling itself linearizable has promised nothing about two objects together. - Paper: "How to Make a Multiprocessor Computer That Correctly Executes Multiprocess Programs", Leslie Lamport, 1979
Two pages, and the origin of sequential consistency: the result of any execution is as if all processors' operations ran in some sequential order, with each processor's own operations in program order. Use for: the definition itself, and for seeing how much weaker it is than linearizability once the real-time clause is absent. - Article: "CAP Twelve Years Later: How the 'Rules' Have Changed", Eric Brewer, InfoQ
CAP's own author revisiting and correcting common misreadings of the theorem, including that partition tolerance isn't optional and that the real trade-off is more nuanced than "pick two". Use for: using CAP correctly instead of the oversimplified version most engineers repeat. - Paper: "In Search of an Understandable Consensus Algorithm", Ongaro and Ousterhout, 2014
The Raft paper, written explicitly to make a consensus protocol's mechanism and cost understandable without requiring a from-scratch proof of correctness. Use for: the primary source on what a consensus protocol actually does and what it costs to run. - Site: "The Raft Consensus Algorithm", raft.github.io
Interactive visualization of Raft's leader election and log replication, letting you watch the protocol handle a simulated node failure or partition. Use for: building intuition for Raft's mechanics before or alongside reading the paper. - Site: "Phenomena", Jepsen
The named anomaly vocabulary: Adya's dependency-based phenomena (G0throughG2), the SQL ones (dirty read, lost update, write skew), the temporal ones (stale read, real-time, process) and long fork. Each page says which models permit the phenomenon and which forbid it. Use for: describing an incident symptom precisely enough to ask whether it was a bug, since whether a phenomenon is legal depends on the model the system claimed. Prefer the Adya family for distributed systems; the SQL phenomena are defined by event order rather than dataflow. -
Site: "Analyses", Jepsen
Real distributed databases and coordination systems tested under actual network partitions and process pauses, with the specific consistency violations each analysis found. Use for: concrete, real-system evidence of what partial failure actually does to a system that assumed the network was reliable. -
Paper: "The Accrual Failure Detector", Hayashibara, Defago, Yared and Katayama, 2004
Introduces accrual failure detection, where the detector reports a suspicion level on a continuous scale instead of a boolean trust-or-suspect, so each application picks its own threshold against a scale the detector adapts to observed network conditions. Thephidetector is the implementation, measured over an intercontinental link. Use for: the detection and timeout mechanics behind telling a slow node from a dead one, and for why a fixed timeout is a bet on a distribution. -
Article: "Dynamo (storage system)", Wikipedia
A summary of Amazon's Dynamo paper (DeCandia et al., 2007) and its techniques table: consistent hashing for partitioning, vector clocks for highly available writes, sloppy quorums and hinted handoff for temporary failures, and Merkle-tree anti-entropy for permanent ones. Use for: leaderless replication as a concrete, real design, and the fact that DynamoDB itself later chose single-leader replication instead, despite the shared name and lineage. - Article: "Quorum (distributed computing)", Wikipedia
Covers Gifford's 1979 quorum-based voting for replicated data: theVr + Vw > Vrule that guarantees a read quorum and a write quorum overlap, and the separateVw > V/2rule that guarantees two write quorums overlap with each other. Use for: the precise arithmetic behindR + W > N, rather than an intuitive but imprecise notion of "majority agreement". - Article: "Consistent hashing", Wikipedia
Covers the ring construction, theO(K/N)average-case bound on keys remapped when a node joins or leaves, and the practical extensions: virtual nodes (to avoid dumping a failed node's whole load onto one neighbor) and replicating a single "hot" key onto multiple contiguous nodes. Use for: precisely why consistent hashing beats plainhash(key) mod M, and the specific gaps a bare ring still leaves that virtual nodes and hot-key replication close. - Article: "Version vector", Wikipedia
Covers how a version vector detects happened-before versus concurrent updates for causality tracking among replicas, and its explicit distinction from a vector clock despite sharing the same underlying state. Use for: the precise mechanism behind detecting whether two writes actually conflict, before any resolution strategy is applied. - Article: "Conflict-free replicated data type", Wikipedia
Covers the state-based (CvRDT) versus operation-based (CmRDT) distinction, the commutative/associative/idempotent properties each requires, and a concrete worked example (the G-Counter, merging by element-wise maximum). Use for: how a CRDT merges concurrent updates without loss, and the delivery-guarantee trade-off between the two CRDT shapes. - Article: "Two-phase commit protocol", Wikipedia
Covers the voting and commit phases, the exact message flow, and the documented blocking failure mode when a coordinator fails after a participant has voted yes. Use for: precisely why 2PC is a blocking protocol, and the specific, worse case where the coordinator and a participant fail together. - Pattern: "Saga", microservices.io
Chris Richardson's pattern reference for sagas: compensating transactions in place of automatic rollback, choreography versus orchestration, and the drawbacks (lost isolation, the dual-write problem each step still faces). Use for: the precise trade-offs a saga makes against 2PC, not just "it's the microservices way to do transactions." - Pattern: "Transactional outbox", microservices.io
Chris Richardson's pattern reference for the outbox table and message relay. Use for: exactly what the pattern guarantees (atomicity between a database commit and a message send, preserved order) and what it explicitly does not (exactly-once delivery, which is why a consumer must be idempotent). - Article: "Exponential Backoff And Jitter", AWS Architecture Blog
Amazon's own worked simulation of capped exponential backoff with and without jitter, under contention from many clients. Use for: the measured effect of jitter (more than halving retry call volume in their 100-client case) versus backoff alone, which still leaves synchronized retry clusters. - Article: "Circuit Breaker", Martin Fowler
The reference description of the circuit breaker pattern (attributed to Michael Nygard's Release It), including the closed/open/half-open state machine. Use for: the precise trip and reset mechanics, and why an open breaker fails fast instead of attempting a call it has already decided is failing. - Specification: "Trace API", OpenTelemetry
The authoritative definition of a span (trace ID, span ID, parent span ID) and the five span kinds. Use for: the precise structure that ties spans into a trace, and the specific, easy-to-miss fact that aPRODUCERspan's duration has no critical-path relationship to itsCONSUMERspan. - Specification: "Context Propagation API", OpenTelemetry
Covers the W3C Trace Contexttraceparent/tracestateheaders and what a propagator must do with them. Use for: exactly what has to be forwarded across a network boundary for a trace to stay connected, and why a missing propagation step silently breaks it in two. - Article: "Analyses", Jepsen
Jepsen's own stated methodology (opaque-box testing of real binaries, testing under actual distributed-systems failure modes, generative testing checked against a formal model), alongside every analysis it has published. Use for: the precise trade-off Jepsen makes, real, production-observable bugs at the cost of nondeterministic, unrepeatable test runs. - Article: "Chaos engineering", Wikipedia
Covers the origin of Chaos Monkey at Netflix, the "Principles of Chaos Engineering" manifesto, and the steady-state-hypothesis framing. Use for: how chaos engineering's aim (does production behavior match what's expected under a specific induced fault) differs from Jepsen's precise, model-checked claim. - Docs: "Simulation and Testing", FoundationDB
FoundationDB's own account of its deterministic simulation testing: full determinism for repeatable debugging, and the time-compression that lets it run the equivalent of roughly a trillion CPU-hours of testing. Use for: the opposite trade from Jepsen and chaos engineering, reproducibility and volume in exchange for testing a simulated system rather than the real production binary.
Gaps
- No source yet on how a real incident is diagnosed end to end, as opposed to the mechanisms individually. Rechecked while writing the stage 5 reference sheet: the Phenomena pages supply the vocabulary for naming a symptom, and the Analyses supply worked examples, but each analysis is written about one system rather than as a method, so the four-question sequence in the sheet is the workspace's own synthesis and rests on no single source. Still open.