Papers in Order

Distributed Systems

From Lamport clocks to globally consistent transactions — the consensus protocols, storage architectures, and coordination primitives that make the cloud actually work.

26 papers 8 levels Included papers 1978 – 2022 MVRP 4 papers
0 of 26 papers read. Progress stays in this browser.

Minimum viable reading path

The 4 papers that give you most of the field's mental model, in reading order.

  1. Time, Clocks, and the Ordering of Events in a Distributed System
  2. The Part-Time Parliament (Paxos)
  3. Bigtable: A Distributed Storage System for Structured Data
  4. Spanner: Google's Globally-Distributed Database
Level 0 Foundations
01 Time, Clocks, and the Ordering of Events in a Distributed System MVRP
TL;DR

Defines logical clocks and the happens-before relation for ordering events without synchronised time.

Why read this

The single most-cited distributed-systems paper. Read it first; everything else assumes it.

Prerequisites

Basic concurrency

Key takeaway

You don't need synchronised clocks — only a partial ordering of events.

Read the paper
02 The Byzantine Generals Problem
TL;DR

Formalises consensus when some nodes can lie or behave arbitrarily, and proves bounds on tolerance.

Why read this

Defines the worst-case adversary model that drives blockchain and high-assurance systems.

Prerequisites

Lamport clocks, basic algorithms

Key takeaway

You can survive arbitrary failures iff fewer than one-third of nodes are faulty.

Read the paper
03 End-to-End Arguments in System Design
TL;DR

Argues that functions like reliability and security belong at endpoints, not inside the network.

Why read this

A foundational systems-design principle, far beyond just distributed systems.

Prerequisites

TCP/IP basics

Key takeaway

Push correctness to the edges; let the middle stay simple.

Read the paper
04 Impossibility of Distributed Consensus with One Faulty Process (FLP)
TL;DR

Proves that no deterministic asynchronous algorithm can guarantee consensus if even one process can crash.

Why read this

The negative result every consensus designer must know — sets the boundary for what's possible.

Prerequisites

Distributed algorithms basics, asynchronous models

Key takeaway

In a fully asynchronous network, you must give up either fault-tolerance or termination.

Read the paper
Level 1 Consensus
05 The Part-Time Parliament (Paxos) MVRP
TL;DR

A consensus protocol that tolerates crashes via majority quorums and ordered ballots.

Why read this

The bedrock consensus protocol — every modern system either runs Paxos or a Paxos cousin.

Prerequisites

FLP, majority voting

Key takeaway

Majorities of replicas can agree on a value despite asynchronous networks and crash failures.

Read the paper
06 In Search of an Understandable Consensus Algorithm (Raft)
TL;DR

A consensus protocol designed for understandability with explicit leader election and log replication.

Why read this

Replaced Paxos in most production systems written after 2014. Easier to read, easier to implement.

Prerequisites

Paxos

Key takeaway

A clearer protocol with the same guarantees is itself a research contribution.

Read the paper
Level 2 CAP & tradeoffs
07 CAP Twelve Years Later: How the "Rules" Have Changed
TL;DR

Brewer revisits his CAP conjecture and clarifies the actual tradeoffs designers face during partitions.

Why read this

Cleaner than the original CAP conjecture; corrects common misreadings.

Prerequisites

Distributed systems basics

Key takeaway

CAP is about what you do during partitions, not a forced binary choice.

Read the paper
08 A Comprehensive Study of CRDTs
TL;DR

A taxonomy of conflict-free replicated data types that converge without coordination.

Why read this

The cleanest reference on the data-type primitives behind eventually-consistent systems and collaborative editors.

Prerequisites

Distributed systems basics, lattice theory

Key takeaway

Pick data types that mathematically can't conflict, and you can drop coordination.

Read the paper
Level 3 Storage
09 The Google File System
TL;DR

A petabyte-scale distributed file system optimised for large appends and commodity hardware failures.

Why read this

The conceptual ancestor of HDFS and the storage layer of an entire generation of big-data stacks.

Prerequisites

File system basics

Key takeaway

Design for failure as the common case, not the exception.

Read the paper
10 Bigtable: A Distributed Storage System for Structured Data MVRP
TL;DR

A sparse, distributed, persistent multi-dimensional sorted map built on GFS and Chubby.

Why read this

The conceptual ancestor of HBase, Cassandra, and most modern wide-column NoSQL stores.

Prerequisites

GFS, LSM-trees

Key takeaway

Sorted-string tables and a log-structured merge engine scale to petabytes.

Read the paper
11 Dynamo: Amazon's Highly Available Key-value Store
TL;DR

A masterless, eventually-consistent KV store using consistent hashing, vector clocks, and quorums.

Why read this

The conceptual root of Cassandra, Riak, and the AP-side of NoSQL.

Prerequisites

Consistent hashing, vector clocks

Key takeaway

Tuneable consistency on writes and reads is a more useful knob than a single consistency level.

Read the paper
12 Cassandra: A Decentralized Structured Storage System
TL;DR

A wide-column store mixing Bigtable's data model with Dynamo's eventually-consistent ring.

Why read this

The pragmatic distillation of two papers everyone reads — Bigtable plus Dynamo.

Prerequisites

Bigtable, Dynamo

Key takeaway

You can fuse two seemingly opposing designs if you pick the right pieces of each.

Read the paper
Level 4 Coordination
13 The Chubby Lock Service for Loosely-Coupled Distributed Systems
TL;DR

A coarse-grained distributed lock service implemented over Paxos for Google's coordination needs.

Why read this

Codifies what a coordination primitive should look like — read this before ZooKeeper.

Prerequisites

Paxos

Key takeaway

A reliable lock service is itself a foundational distributed primitive.

Read the paper
14 ZooKeeper: Wait-free Coordination for Internet-scale Systems
TL;DR

A high-throughput coordination service exposing a hierarchical filesystem-like API over a replicated log.

Why read this

The most-deployed coordination service of the last decade; understanding it sharpens your distributed-systems intuition.

Prerequisites

Chubby, Paxos

Key takeaway

Expose coordination through familiar primitives (files, watches) and people will actually use it.

Read the paper
Level 5 Computation
15 MapReduce: Simplified Data Processing on Large Clusters
TL;DR

A two-stage programming model with automatic parallelisation, fault tolerance, and data locality.

Why read this

The paper that defined "big data" and seeded Hadoop, Spark, and the modern analytics stack.

Prerequisites

Functional programming basics

Key takeaway

Reduce a complex distributed problem to two simple functions, and the framework handles the rest.

Read the paper
16 Resilient Distributed Datasets: A Fault-Tolerant Abstraction (Spark)
TL;DR

An in-memory abstraction that recomputes lost partitions from lineage rather than checkpointing.

Why read this

The paper that displaced MapReduce; required reading for anyone touching analytics infrastructure.

Prerequisites

MapReduce, functional programming

Key takeaway

Logging deterministic transformations is enough to reconstruct lost state.

Read the paper
17 Apache Flink: Stream and Batch Processing in a Single Engine
TL;DR

A unified runtime for streaming and batch via consistent checkpointing and asynchronous barrier snapshots.

Why read this

Best modern reference on what stream processing actually requires.

Prerequisites

Spark, basic stream processing

Key takeaway

Batch is just a special case of streaming with bounded data.

Read the paper
18 Mesos: A Platform for Fine-Grained Resource Sharing
TL;DR

A two-level scheduler that offers resources to higher-level frameworks rather than scheduling tasks itself.

Why read this

The conceptual ancestor of Kubernetes; clarifies the design space of cluster schedulers.

Prerequisites

Distributed systems basics

Key takeaway

Push policy decisions up to the framework; keep the cluster manager mechanism-only.

Read the paper
Level 6 Globally consistent
19 Spanner: Google's Globally-Distributed Database MVRP
TL;DR

A globally-distributed SQL database providing externally-consistent transactions via TrueTime.

Why read this

The first system to make globally-consistent transactions practical — a watershed in DB design.

Prerequisites

Bigtable, Paxos, clock synchronisation

Key takeaway

Tight clock bounds, plus a wait, give you serialisable distributed transactions.

Read the paper
20 Calvin: Fast Distributed Transactions for Partitioned Database Systems
TL;DR

A deterministic database that orders all transactions globally before executing them in parallel.

Why read this

A fundamentally different approach to transactions than 2PC — worth reading for the contrast.

Prerequisites

Spanner, 2PC

Key takeaway

If you agree on the order of transactions first, you can run them without coordinating during execution.

Read the paper
Level 7 Production systems
21 Kafka: A Distributed Messaging System for Log Processing
TL;DR

A distributed, partitioned, replicated commit log designed as the central nervous system of an org's data.

Why read this

The paper that re-popularised the log as a foundational data abstraction.

Prerequisites

Distributed systems basics

Key takeaway

A simple, durable, ordered append-only log is the right primitive for moving data around.

Read the paper
22 Large-scale Cluster Management at Google with Borg
TL;DR

A decade of lessons from running Google's internal cluster manager — the predecessor to Kubernetes.

Why read this

The most candid production-systems paper of the 2010s; required reading before designing any scheduler.

Prerequisites

Mesos

Key takeaway

Most of the work in cluster management is operational, not algorithmic.

Read the paper
23 Dapper, a Large-Scale Distributed Systems Tracing Infrastructure
TL;DR

A low-overhead, always-on tracing system that propagates span identifiers across RPC boundaries.

Why read this

Sets the design template for OpenTelemetry, Jaeger, and every modern tracing system.

Prerequisites

RPC basics

Key takeaway

Sample your traces, propagate context, and you can debug anything.

Read the paper
24 Serverless Computing: One Step Forward, Two Steps Back
TL;DR

A measured critique of FaaS that catalogues what serverless gets right and where it falls short.

Why read this

A useful counterweight to vendor hype; sharpens your sense of what cloud abstractions should do.

Prerequisites

Distributed systems basics, cloud architecture

Key takeaway

Stateless functions are easy; the hard problems are state, coordination, and shared memory.

Read the paper
25 Anna: A KVS for Any Scale
TL;DR

A coordination-free, multi-master KV store that scales linearly across cores, machines, and regions.

Why read this

A clean modern read on what coordination-free design buys you and where it bites.

Prerequisites

Dynamo, CRDTs

Key takeaway

Lattice-based replication can give you both speed and consistency knobs.

Read the paper
26 TLA+: Specifying Systems
TL;DR

A specification language and model checker for concurrent and distributed algorithms, with worked examples.

Why read this

Less a paper than a discipline — distributed bugs become tractable when you can specify them.

Prerequisites

Discrete math, first-order logic

Key takeaway

Specifying a system precisely catches bugs no test ever will.

Read the paper