Part 9 · 1 chapters · ~12 min
The Papers
The papers that defined distributed systems on one timeline, from Lamport's clocks through FLP, Paxos, CAP, Google's systems, Dynamo, ZooKeeper, Kafka, Spanner and Calvin to Raft, each with what it proved, what changed, and where it appears in this course.
18
What each paper proved
The papers below are worth reading in full. For each one: what it proved or built, what changed because of it, and where it shows up in this course.
| paper | year | what it proved or built | what changed | this course |
|---|---|---|---|---|
| Time, Clocks, and the Ordering of Events (Lamport) | 1978 | happened-before; logical clocks; state machine replication | ordering without physical time became the foundation | part 1 |
| The Byzantine Generals Problem (Lamport, Shostak, Pease) | 1982 | agreement needs 3f+1 nodes to tolerate f arbitrary faults | BFT protocols; later, blockchains | part 2 (context) |
| Impossibility of Distributed Consensus with One Faulty Process (FLP) | 1985 | no deterministic consensus guaranteed in pure asynchrony with one crash | protocols rely on timeouts and partial synchrony for liveness | parts 2, 6 |
| The Part-Time Parliament / Paxos Made Simple (Lamport) | 1998 / 2001 | consensus under crash failures with majorities | Chubby, Spanner, many databases | part 2 |
| Brewer's conjecture and its proof (Gilbert, Lynch) | 2000 / 2002 | consistency or availability during a partition | the vocabulary of every database trade-off | part 5 |
| The Google File System | 2003 | large files on commodity disks with a single master and chunk replication | HDFS and the big-data era | parts 3, 4 |
| MapReduce | 2004 | batch computation over thousands of machines with automatic retry | Hadoop; later Spark and dataflow | part 7 (context) |
| Bigtable | 2006 | a sparse, sorted, range-partitioned table at petabyte scale | HBase, Cassandra's data model, Cloud Bigtable | part 4 |
| The Chubby Lock Service | 2006 | a Paxos-backed lock service with leases for coarse coordination | ZooKeeper, etcd as coordination services | parts 2, 6 |
| Dynamo | 2007 | always-writable leaderless store: consistent hashing, sloppy quorums, vector clocks | Cassandra, Riak, DynamoDB's lineage | parts 1, 3, 4 |
| ZooKeeper | 2010 | wait-free coordination primitives on an ordered, replicated log (ZAB) | Kafka (pre-KRaft), Hadoop, HBase coordination | part 2 |
| Kafka | 2011 | a partitioned, replicated log as a messaging and data backbone | event-driven architectures, CDC, stream processing | part 7 |
| Spanner | 2012 | externally consistent global transactions with TrueTime | CockroachDB, YugabyteDB; bounded-clock designs | parts 1, 5 |
| Calvin | 2012 | agree on transaction order first, execute deterministically: no 2PC | FaunaDB; deterministic database research | part 8 |
| Raft | 2014 | understandable consensus equivalent to multi-Paxos | etcd, Consul, CockroachDB, TiKV, KRaft | part 2 |
how to read them
Read the introduction and conclusion first, then the system design section, and skip the proofs on a first pass. For each paper, write one sentence on what it made possible and one on what it gave up. Martin Kleppmann's Designing Data-Intensive Applications is the best companion: it cites nearly all of these and explains them in context.
THE PAPERS, ON A TIMELINE
the results that defined what is possible, and the systems that proved it at scale
swipe the figure sideways, or tap expand for full screen
1/6
foundations
Foundations: Lamport, "Time, Clocks, and the Ordering of Events in a Distributed System" (1978) defines happened-before and logical clocks (part 1). Lamport, Shostak and Pease, "The Byzantine Generals Problem" (1982) defines agreement with arbitrarily faulty nodes, the basis of BFT and blockchains.