Distributed Systems
From Lamport clocks to globally consistent transactions — the consensus protocols, storage architectures, and coordination primitives that make the cloud actually work.
Minimum viable reading path
The 4 papers that give you most of the field's mental model, in reading order.
01 Time, Clocks, and the Ordering of Events in a Distributed System MVRP
Defines logical clocks and the happens-before relation for ordering events without synchronised time.
The single most-cited distributed-systems paper. Read it first; everything else assumes it.
Basic concurrency
You don't need synchronised clocks — only a partial ordering of events.
02 The Byzantine Generals Problem
Formalises consensus when some nodes can lie or behave arbitrarily, and proves bounds on tolerance.
Defines the worst-case adversary model that drives blockchain and high-assurance systems.
Lamport clocks, basic algorithms
You can survive arbitrary failures iff fewer than one-third of nodes are faulty.
03 End-to-End Arguments in System Design
Argues that functions like reliability and security belong at endpoints, not inside the network.
A foundational systems-design principle, far beyond just distributed systems.
TCP/IP basics
Push correctness to the edges; let the middle stay simple.
04 Impossibility of Distributed Consensus with One Faulty Process (FLP)
Proves that no deterministic asynchronous algorithm can guarantee consensus if even one process can crash.
The negative result every consensus designer must know — sets the boundary for what's possible.
Distributed algorithms basics, asynchronous models
In a fully asynchronous network, you must give up either fault-tolerance or termination.
05 The Part-Time Parliament (Paxos) MVRP
A consensus protocol that tolerates crashes via majority quorums and ordered ballots.
The bedrock consensus protocol — every modern system either runs Paxos or a Paxos cousin.
FLP, majority voting
Majorities of replicas can agree on a value despite asynchronous networks and crash failures.
06 In Search of an Understandable Consensus Algorithm (Raft)
A consensus protocol designed for understandability with explicit leader election and log replication.
Replaced Paxos in most production systems written after 2014. Easier to read, easier to implement.
Paxos
A clearer protocol with the same guarantees is itself a research contribution.
07 CAP Twelve Years Later: How the "Rules" Have Changed
Brewer revisits his CAP conjecture and clarifies the actual tradeoffs designers face during partitions.
Cleaner than the original CAP conjecture; corrects common misreadings.
Distributed systems basics
CAP is about what you do during partitions, not a forced binary choice.
08 A Comprehensive Study of CRDTs
A taxonomy of conflict-free replicated data types that converge without coordination.
The cleanest reference on the data-type primitives behind eventually-consistent systems and collaborative editors.
Distributed systems basics, lattice theory
Pick data types that mathematically can't conflict, and you can drop coordination.
09 The Google File System
A petabyte-scale distributed file system optimised for large appends and commodity hardware failures.
The conceptual ancestor of HDFS and the storage layer of an entire generation of big-data stacks.
File system basics
Design for failure as the common case, not the exception.
10 Bigtable: A Distributed Storage System for Structured Data MVRP
A sparse, distributed, persistent multi-dimensional sorted map built on GFS and Chubby.
The conceptual ancestor of HBase, Cassandra, and most modern wide-column NoSQL stores.
GFS, LSM-trees
Sorted-string tables and a log-structured merge engine scale to petabytes.
11 Dynamo: Amazon's Highly Available Key-value Store
A masterless, eventually-consistent KV store using consistent hashing, vector clocks, and quorums.
The conceptual root of Cassandra, Riak, and the AP-side of NoSQL.
Consistent hashing, vector clocks
Tuneable consistency on writes and reads is a more useful knob than a single consistency level.
12 Cassandra: A Decentralized Structured Storage System
A wide-column store mixing Bigtable's data model with Dynamo's eventually-consistent ring.
The pragmatic distillation of two papers everyone reads — Bigtable plus Dynamo.
Bigtable, Dynamo
You can fuse two seemingly opposing designs if you pick the right pieces of each.
13 The Chubby Lock Service for Loosely-Coupled Distributed Systems
A coarse-grained distributed lock service implemented over Paxos for Google's coordination needs.
Codifies what a coordination primitive should look like — read this before ZooKeeper.
Paxos
A reliable lock service is itself a foundational distributed primitive.
14 ZooKeeper: Wait-free Coordination for Internet-scale Systems
A high-throughput coordination service exposing a hierarchical filesystem-like API over a replicated log.
The most-deployed coordination service of the last decade; understanding it sharpens your distributed-systems intuition.
Chubby, Paxos
Expose coordination through familiar primitives (files, watches) and people will actually use it.
15 MapReduce: Simplified Data Processing on Large Clusters
A two-stage programming model with automatic parallelisation, fault tolerance, and data locality.
The paper that defined "big data" and seeded Hadoop, Spark, and the modern analytics stack.
Functional programming basics
Reduce a complex distributed problem to two simple functions, and the framework handles the rest.
16 Resilient Distributed Datasets: A Fault-Tolerant Abstraction (Spark)
An in-memory abstraction that recomputes lost partitions from lineage rather than checkpointing.
The paper that displaced MapReduce; required reading for anyone touching analytics infrastructure.
MapReduce, functional programming
Logging deterministic transformations is enough to reconstruct lost state.
17 Apache Flink: Stream and Batch Processing in a Single Engine
A unified runtime for streaming and batch via consistent checkpointing and asynchronous barrier snapshots.
Best modern reference on what stream processing actually requires.
Spark, basic stream processing
Batch is just a special case of streaming with bounded data.
18 Mesos: A Platform for Fine-Grained Resource Sharing
A two-level scheduler that offers resources to higher-level frameworks rather than scheduling tasks itself.
The conceptual ancestor of Kubernetes; clarifies the design space of cluster schedulers.
Distributed systems basics
Push policy decisions up to the framework; keep the cluster manager mechanism-only.
19 Spanner: Google's Globally-Distributed Database MVRP
A globally-distributed SQL database providing externally-consistent transactions via TrueTime.
The first system to make globally-consistent transactions practical — a watershed in DB design.
Bigtable, Paxos, clock synchronisation
Tight clock bounds, plus a wait, give you serialisable distributed transactions.
20 Calvin: Fast Distributed Transactions for Partitioned Database Systems
A deterministic database that orders all transactions globally before executing them in parallel.
A fundamentally different approach to transactions than 2PC — worth reading for the contrast.
Spanner, 2PC
If you agree on the order of transactions first, you can run them without coordinating during execution.
21 Kafka: A Distributed Messaging System for Log Processing
A distributed, partitioned, replicated commit log designed as the central nervous system of an org's data.
The paper that re-popularised the log as a foundational data abstraction.
Distributed systems basics
A simple, durable, ordered append-only log is the right primitive for moving data around.
22 Large-scale Cluster Management at Google with Borg
A decade of lessons from running Google's internal cluster manager — the predecessor to Kubernetes.
The most candid production-systems paper of the 2010s; required reading before designing any scheduler.
Mesos
Most of the work in cluster management is operational, not algorithmic.
23 Dapper, a Large-Scale Distributed Systems Tracing Infrastructure
A low-overhead, always-on tracing system that propagates span identifiers across RPC boundaries.
Sets the design template for OpenTelemetry, Jaeger, and every modern tracing system.
RPC basics
Sample your traces, propagate context, and you can debug anything.
24 Serverless Computing: One Step Forward, Two Steps Back
A measured critique of FaaS that catalogues what serverless gets right and where it falls short.
A useful counterweight to vendor hype; sharpens your sense of what cloud abstractions should do.
Distributed systems basics, cloud architecture
Stateless functions are easy; the hard problems are state, coordination, and shared memory.
25 Anna: A KVS for Any Scale
A coordination-free, multi-master KV store that scales linearly across cores, machines, and regions.
A clean modern read on what coordination-free design buys you and where it bites.
Dynamo, CRDTs
Lattice-based replication can give you both speed and consistency knobs.
26 TLA+: Specifying Systems
A specification language and model checker for concurrent and distributed algorithms, with worked examples.
Less a paper than a discipline — distributed bugs become tractable when you can specify them.
Discrete math, first-order logic
Specifying a system precisely catches bugs no test ever will.