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.

paperyearwhat it proved or builtwhat changedthis course
Time, Clocks, and the Ordering of Events (Lamport)1978happened-before; logical clocks; state machine replicationordering without physical time became the foundationpart 1
The Byzantine Generals Problem (Lamport, Shostak, Pease)1982agreement needs 3f+1 nodes to tolerate f arbitrary faultsBFT protocols; later, blockchainspart 2 (context)
Impossibility of Distributed Consensus with One Faulty Process (FLP)1985no deterministic consensus guaranteed in pure asynchrony with one crashprotocols rely on timeouts and partial synchrony for livenessparts 2, 6
The Part-Time Parliament / Paxos Made Simple (Lamport)1998 / 2001consensus under crash failures with majoritiesChubby, Spanner, many databasespart 2
Brewer's conjecture and its proof (Gilbert, Lynch)2000 / 2002consistency or availability during a partitionthe vocabulary of every database trade-offpart 5
The Google File System2003large files on commodity disks with a single master and chunk replicationHDFS and the big-data eraparts 3, 4
MapReduce2004batch computation over thousands of machines with automatic retryHadoop; later Spark and dataflowpart 7 (context)
Bigtable2006a sparse, sorted, range-partitioned table at petabyte scaleHBase, Cassandra's data model, Cloud Bigtablepart 4
The Chubby Lock Service2006a Paxos-backed lock service with leases for coarse coordinationZooKeeper, etcd as coordination servicesparts 2, 6
Dynamo2007always-writable leaderless store: consistent hashing, sloppy quorums, vector clocksCassandra, Riak, DynamoDB's lineageparts 1, 3, 4
ZooKeeper2010wait-free coordination primitives on an ordered, replicated log (ZAB)Kafka (pre-KRaft), Hadoop, HBase coordinationpart 2
Kafka2011a partitioned, replicated log as a messaging and data backboneevent-driven architectures, CDC, stream processingpart 7
Spanner2012externally consistent global transactions with TrueTimeCockroachDB, YugabyteDB; bounded-clock designsparts 1, 5
Calvin2012agree on transaction order first, execute deterministically: no 2PCFaunaDB; deterministic database researchpart 8
Raft2014understandable consensus equivalent to multi-Paxosetcd, Consul, CockroachDB, TiKV, KRaftpart 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.