Distributed Systems

Replication, consistency, coordination, consensus, fault tolerance, partitioning and the trade-offs behind systems that span multiple machines.

Distributed systems are one of the areas where returning to theory changed the way I look at technologies I already use in practice.

Databases, message brokers, microservices, cloud platforms and blockchains all distribute computation or state across multiple machines. What makes them difficult is not simply the number of nodes involved, but the fact that communication can fail, clocks disagree, machines crash and different participants can observe different versions of reality.

How do several machines agree on a value? What happens when part of the network becomes unreachable? When can replicated data be considered consistent? How many failures can a system tolerate? What does “eventual consistency” actually mean for an application?

This section combines two complementary perspectives: the theoretical foundations of distributed computing and the architectural trade-offs that appear in real data-intensive systems.


Topics in This Section

Distributed Architectures · Communication · RPC · Naming · Logical Clocks · Coordination · Leader Election · Consensus · Replication · Consistency Models · CAP · Network Partitions · Fault Tolerance · Crash Faults · Byzantine Faults · Paxos · Distributed Transactions · Sharding · Eventual Consistency · Batch Processing · Stream Processing


Distributed Systems

Maarten van Steen & Andrew S. Tanenbaum

4th Edition — distributed-systems.net

Level
Foundation → Advanced

Best for
Distributed-systems theory, coordination, consistency, replication and fault tolerance.

Distributed Systems provides the theoretical framework I use to understand what changes when computation moves from a single machine to a collection of independent nodes.

What I particularly value is the progression from architectural models and communication to coordination, replication and failures. These concepts make it much easier to reason about real distributed technologies without tying the explanation to one particular product.

The fourth edition also integrates modern examples, including blockchain systems, while retaining the underlying principles that apply across distributed architectures.

What I Use It For

  • understanding what makes a system distributed;
  • architectural models and middleware;
  • distributed processes and communication;
  • Remote Procedure Calls and message-based communication;
  • naming and resource discovery;
  • logical clocks and event ordering;
  • distributed coordination;
  • leader election and mutual exclusion;
  • replication strategies;
  • consistency models;
  • fault models;
  • crash fault tolerance;
  • Byzantine fault tolerance;
  • consensus and Paxos;
  • distributed commit and recovery.

Chapters Worth Reading

Foundations & Architectures

Chapter 1 — Introduction
Characteristics, goals, challenges and fundamental concepts of distributed systems.

Chapter 2 — Architectures
Architectural styles, middleware and different ways of organising distributed components.

These chapters provide the mental model for understanding why distribution introduces scalability opportunities but also communication delays, partial failures and coordination problems.

Processes & Communication

Chapter 3 — Processes
Processes, threads, clients, servers and the execution models used in distributed environments.

Chapter 4 — Communication
Communication between distributed components, including message-oriented approaches and Remote Procedure Calls.

Naming

Chapter 5 — Naming
How distributed systems identify and locate entities, resources and services.

This chapter is useful for connecting abstract naming concepts with practical mechanisms such as hierarchical naming and distributed lookup systems.

Coordination

Chapter 6 — Coordination
How independent nodes coordinate actions when there is no single globally shared clock or state.

  • physical and logical clocks;
  • ordering of events;
  • distributed mutual exclusion;
  • leader election;
  • coordination between processes.
Consistency & Replication

Chapter 7 — Consistency and Replication
How replicated systems balance performance, availability and the guarantees offered to clients.

  • replication strategies;
  • data-centric consistency models;
  • client-centric consistency;
  • eventual consistency;
  • replicated-write protocols;
  • cache coherence.
Fault Tolerance & Consensus

Chapter 8 — Fault Tolerance
How distributed systems continue operating when processes, machines or communications fail.

  • failure models;
  • redundancy;
  • process resilience;
  • crash-fault consensus;
  • Paxos;
  • arbitrary and Byzantine failures;
  • consensus in blockchain systems;
  • failure detection;
  • reliable client-server communication;
  • distributed commit;
  • checkpointing and recovery.
Distributed-System Security

Chapter 9 — Security
Security requirements and mechanisms in systems where components communicate across potentially untrusted networks.

The chapter connects distributed-system security with cryptography, authentication, access control and secure communication.

My Suggested Learning Path

Distributed-System Foundations
Chapters 1–2

Processes & Communication
Chapters 3–4

Coordination
Chapters 5–6

Consistency & Replication
Chapter 7

Failures & Consensus
Chapter 8

Security
Chapter 9


Designing Data-Intensive Applications

Martin Kleppmann & Chris Riccomini

2nd Edition — O’Reilly Media

Level
Intermediate → Advanced

Best for
Connecting distributed-systems theory with the architecture of real databases, data platforms and cloud applications.

Designing Data-Intensive Applications approaches distributed systems from a different direction.

Rather than starting with distributed algorithms, it starts with the design decisions engineers face when building systems that need to remain reliable, scalable and maintainable while storing and processing large amounts of data.

I find this perspective especially valuable because it makes abstract concepts such as replication, consistency, partitioning and consensus immediately relevant to architectural decisions.

What I Use It For

  • reasoning about scalability and reliability;
  • understanding architectural trade-offs;
  • replication strategies;
  • sharding and partitioning;
  • distributed transactions;
  • consistency and isolation;
  • partial failures and unreliable networks;
  • logical clocks;
  • consensus;
  • event-driven architecture;
  • batch processing;
  • stream processing;
  • designing data flows across multiple systems.

Chapters Worth Reading

Architecture, Reliability & Scalability

Chapter 1 — Trade-Offs in Data Systems Architecture
Operational versus analytical systems, cloud versus self-hosting, distributed versus single-node architectures, microservices and serverless systems.

Chapter 2 — Defining Nonfunctional Requirements
Performance, latency, reliability, fault tolerance, scalability, operability and maintainability.

Communication & Evolution

Chapter 5 — Encoding and Evolution
Data formats, schema evolution and the different ways data flows between databases, services and event-driven systems.

  • JSON and binary formats;
  • Protocol Buffers;
  • Avro;
  • REST and RPC;
  • event-driven architectures;
  • durable workflows.
Replication

Chapter 6 — Replication
How multiple copies of data are maintained across different nodes and regions.

  • single-leader replication;
  • synchronous and asynchronous replication;
  • replication lag;
  • multi-leader replication;
  • leaderless replication;
  • conflicting writes;
  • multi-region operation.
Sharding

Chapter 7 — Sharding
Partitioning data across multiple nodes to distribute storage and workload.

  • key-range sharding;
  • hash-based sharding;
  • hot spots;
  • rebalancing;
  • request routing;
  • secondary indexes across shards.
Transactions

Chapter 8 — Transactions
Transaction guarantees, isolation anomalies, serializability and transactions that span distributed components.

  • ACID;
  • Read Committed;
  • Snapshot Isolation;
  • lost updates;
  • write skew and phantoms;
  • serializability;
  • Two-Phase Locking;
  • Serializable Snapshot Isolation;
  • distributed transactions;
  • Two-Phase Commit.
Failures in Distributed Systems

Chapter 9 — The Trouble with Distributed Systems
Why failures in distributed environments differ fundamentally from failures inside a single machine.

  • partial failures;
  • unreliable networks;
  • timeouts;
  • unbounded delays;
  • unreliable clocks;
  • process pauses;
  • distributed locks and leases;
  • Byzantine faults;
  • system models.
Consistency & Consensus

Chapter 10 — Consistency and Consensus
The guarantees distributed systems can provide when multiple nodes need to agree on shared state.

  • linearizability;
  • logical clocks;
  • distributed ID generation;
  • consensus;
  • consensus algorithms in practice;
  • coordination services.
Batch & Stream Processing

Chapter 11 — Batch Processing
Distributed filesystems, MapReduce, dataflow engines, ETL and distributed processing of bounded datasets.

Chapter 12 — Stream Processing
Event streams, message brokers, change-data capture, stream joins, time and fault tolerance.

Chapter 13 — A Philosophy of Streaming Systems
Data integration, derived state, dataflow architectures and reasoning about correctness across multiple systems.

My Suggested Learning Path

Architecture & Scalability
Chapters 1–2

Replication & Sharding
Chapters 6–7

Transactions
Chapter 8

Failures & Consensus
Chapters 9–10

Distributed Data Processing
Chapters 11–13

Communication & Evolution
Chapter 5


Why I Keep Both Books

Van Steen & Tanenbaum

Principles first.

  • coordination;
  • logical clocks;
  • consistency models;
  • replication;
  • failure models;
  • consensus;
  • fault tolerance.

Kleppmann & Riccomini

Trade-offs first.

  • real data architectures;
  • replication strategies;
  • sharding;
  • distributed transactions;
  • consistency in practice;
  • batch and streaming systems;
  • architectural decision-making.

One book helps me understand the principles of distributed computing. The other helps me recognise those principles inside real systems.


Topic → Book Map

Distributed-System Foundations

Distributed Systems: Chapters 1–2
DDIA: Chapters 1–2

Communication

Distributed Systems: Chapter 4
DDIA: Chapter 5

Coordination & Logical Clocks

Distributed Systems: Chapter 6
DDIA: Chapters 9–10

Replication & Consistency

Distributed Systems: Chapter 7
DDIA: Chapters 6 and 10

Partitioning & Sharding

Distributed Systems: addressed across architecture and replication topics
DDIA: Chapter 7

Fault Tolerance

Distributed Systems: Chapter 8
DDIA: Chapter 9

Consensus

Distributed Systems: Chapter 8
DDIA: Chapter 10

Distributed Transactions

Distributed Systems: Chapter 8
DDIA: Chapter 8

Batch & Stream Processing

DDIA: Chapters 11–13


How I Use These Books

I find distributed systems easier to understand when I begin with a failure or design requirement and then identify the principle behind it.

“I have three replicas. One becomes unreachable. Can the system continue accepting writes?”

That question leads to replication, quorum decisions, consistency requirements and the distinction between availability and correctness.

“Two nodes disagree about which event happened first.”

That leads to physical clocks, logical clocks, causality and distributed ordering.

“A network request times out. Did the remote operation fail, or did only the response fail to reach me?”

That leads directly to partial failures, retry semantics, idempotency and the fundamental uncertainty introduced by distributed communication.

“Several nodes need to agree on the same value even when some of them fail.”

That points toward consensus, quorum-based protocols, Paxos and the assumptions required to tolerate crash or Byzantine faults.

Start from the failure or invariant, identify the distributed-systems principle behind it, then study the protocol used to preserve that property.


Related Areas

Distributed systems sit at the intersection of many other areas in this library.

Databases & Data Management

Replication, distributed transactions, consistency, isolation and partitioning.

Software Architecture & Microservices

Service communication, eventual consistency, Saga, messaging and failure handling.

Computer Networks

Latency, unreliable communication, network partitions and message delivery.

Containers, Cloud & DevOps

Distributed deployment, scheduling, replication, autoscaling and failure recovery.

Blockchain & Smart Contracts

Consensus, replicated state, Byzantine faults and distributed agreement.

Cybersecurity & Cryptography

Authentication, trust, adversarial participants and secure distributed communication.


← Back to My Technical Library