Patterns of Distributed Systems Link to heading

Based on: Catalog of Patterns of Distributed Systems by Unmesh Joshi (published on martinfowler.com, Nov 2023)


Table of Contents Link to heading

  1. Quick Reference Summary Table

1. Quick Reference Summary Table Link to heading

# Pattern Category One-Line Problem Core Idea Key Systems
1 Write-Ahead Log Durability Survive crashes without losing state Append every change to a sequential log before applying it PostgreSQL WAL, Kafka, RocksDB
2 Segmented Log Durability Single log file grows unbounded Split log into fixed-size segment files; delete old segments Kafka, etcd, Bookkeeper
3 Low-Water Mark Durability Can’t tell which old log entries are safe to delete Track the minimum log index confirmed by all followers Kafka, Raft implementations
4 High-Water Mark Durability Followers might expose uncommitted entries to clients Track the highest log index replicated to a majority Kafka, Raft, ZooKeeper
5 Versioned Value Durability Overwriting loses history; concurrent writes conflict Store every update with a new version number etcd, CockroachDB, Cassandra
6 HeartBeat Coordination Can’t tell if a remote node is alive or crashed Periodic messages; silence = failure ZooKeeper, Kubernetes, Raft
7 Leader and Followers Coordination Multiple writers cause conflicting state Single leader serializes all writes; followers replicate MySQL, MongoDB, Raft
8 Generation Clock Coordination Old leader comes back and causes split-brain Monotonic epoch counter; reject messages from stale epochs Raft (term), Paxos (ballot)
9 Emergent Leader Coordination Explicit elections are complex Oldest node in cluster is automatically the leader Akka Cluster
10 Lease Coordination Crash-proof distributed locks needed Time-limited lock; must be renewed; expires on crash etcd, ZooKeeper, Chubby
11 Majority Quorum Consensus All-or-nothing availability vs. data loss trade-off Require ⌊N/2⌋+1 confirmations; any two majorities overlap Raft, Paxos, ZooKeeper
12 Replicated Log Consensus Single-node WAL isn’t fault tolerant Ship WAL entries to followers; commit on majority ack Raft (etcd), ZooKeeper ZAB
13 Paxos Consensus Nodes must agree on a single value despite failures Two-phase protocol with ballot numbers to prevent conflicting commits Chubby, ZooKeeper (inspired)
14 Consistent Core Consensus Running consensus on every node in a large cluster is expensive Small 3–5 node strongly-consistent cluster as coordination oracle ZooKeeper, etcd, Consul
15 Lamport Clock Ordering Wall-clock timestamps are unreliable across nodes Logical counter; max(local, received) + 1 on message receipt Distributed version numbers
16 Hybrid Clock Ordering Lamport clocks aren’t human-readable; wall clocks aren’t safe Combine physical time + logical counter CockroachDB, MongoDB, YugabyteDB
17 Clock-Bound Wait Ordering Clock uncertainty causes ordering violations in global transactions Wait before committing so all clocks have passed the same point Google Spanner, CockroachDB
18 Version Vector Ordering Need to detect concurrent writes in multi-master systems Per-node counter vector; concurrent if neither vector dominates Riak, Amazon Dynamo, CouchDB
19 Fixed Partitions Partitioning Data too large for one node; need stable routing Fixed bucket count (e.g., 1024); reassign buckets, not data Kafka, Cassandra vnodes, Redis Cluster
20 Key-Range Partitions Partitioning Range queries require scanning all nodes Co-locate sorted key ranges on nodes; split/merge as needed Bigtable, HBase, CockroachDB
21 Single-Socket Channel Networking Messages arrive out of order with multiple connections One TCP connection per peer pair; TCP guarantees order ZooKeeper, Raft implementations
22 Request Pipeline Networking Waiting for each response limits throughput Send next request before receiving previous response Redis pipeline, HTTP/2, Raft replication
23 Request Batch Networking Per-request network overhead dominates processing cost Accumulate multiple requests; send as one payload Kafka producer, Raft heartbeat
24 Singular Update Queue Concurrency Locks and concurrent state mutation cause bugs Single thread processes all state changes from a queue etcd, ZooKeeper, Node.js event loop
25 Request Waiting List Concurrency Server must track which client waits for which cluster response Map of in-flight requests → callbacks triggered on quorum ack Raft leaders, ZooKeeper
26 Idempotent Receiver Reliability Client retries on timeout cause duplicate side effects Unique request ID; server deduplicates using cached results Stripe, AWS APIs, Kafka EOS
27 Gossip Dissemination Propagation Broadcasting cluster state to all N nodes is O(N) bottleneck Each node randomly picks peers to exchange state; spreads exponentially Cassandra, Dynamo, Consul, Riak
28 Follower Reads Scalability Leader is bottleneck for read-heavy workloads Serve reads from followers (with bounded staleness) MySQL replicas, Cassandra, MongoDB
29 State Watch Observation Polling for cluster state changes wastes resources Server pushes notifications to registered watchers on key changes ZooKeeper, etcd, Kubernetes
30 Two-Phase Commit Transactions Atomic updates across multiple nodes are not naturally safe Coordinator collects votes (prepare), then issues global commit/abort XA/JTA, CockroachDB, Spanner

Pattern Selection Guide Link to heading

If you need… Consider these patterns
Crash recovery for a single node Write-Ahead Log, Segmented Log, Low-Water Mark
Fault-tolerant data replication Replicated Log, Leader and Followers, Majority Quorum
Leader election Leader and Followers, Generation Clock, Emergent Leader, Consistent Core
Ordering events across nodes Lamport Clock, Hybrid Clock, Version Vector
Handling clock skew in global transactions Clock-Bound Wait, Hybrid Clock
Conflict detection in multi-master systems Version Vector, Versioned Value
Scaling read throughput Follower Reads
Efficient cluster information spread Gossip Dissemination
Reliable distributed locking Lease, Consistent Core
Safe client retries Idempotent Receiver
Atomic multi-node updates Two-Phase Commit
Reactive updates to cluster state State Watch
High-throughput messaging Request Pipeline, Request Batch, Single-Socket Channel
Concurrent state management Singular Update Queue, Request Waiting List
Distributing data across nodes Fixed Partitions, Key-Range Partitions