Chapter 02 · Distributed Systems Foundations: Nodes, Networks, Failure, and State
Processes, Nodes, Clusters, Partitions, Replicas, Coordinators, and Client Routing
Trace an AtlasMart request across routing metadata, partitions, coordinators, replicas, acknowledgements, write-ahead-log state, and repair while defining the core vocabulary of a distributed database.
Learning outcomes
Chapter 01 gave AtlasMart a decision notebook but deliberately avoided choosing a distributed database. Chapter 02 now supplies the substrate that every later design depends on. Imagine a checkout request entering a system with six database processes, four data partitions, three copies of each partition, and routing metadata that can change while the request is in flight. A user still sees one Place order button, but the implementation may involve several independent processes and network messages before the operation can be acknowledged.
Define process, node, cluster, shard/partition, replica, coordinator, leader, follower, client, proxy, and routing metadata without treating them as interchangeable words.
Trace one AtlasMart key from client routing through a coordinator, replica acknowledgements, persisted write-ahead-log evidence, and later repair.
Distinguish partitioning, which assigns different data ranges/keys, from replication, which creates additional copies of the same logical data.
Explain why routing metadata needs a version/epoch and why clients cannot safely hard-code one server forever.
Interpret acknowledgement and replica state as evidence while avoiding claims that an acknowledgement count alone proves a particular consistency or durability guarantee.
Chapter 01 identified orders/payments/inventory as correctness-sensitive, catalog/search as locality- and retrieval-sensitive, and sessions/idempotency as key-oriented. Those workload labels now become routes through independent nodes. The same business key can be partitioned, replicated, temporarily divergent, and observed differently by different clients.
The mandatory lab is a deterministic Python standard-library simulation. Python 3.14.7 is the current stable Python feature release at review time; the generation environment executed the lab with Python 3.13.5 on Linux. There is no database product, cloud account, container, external network, or paid feature in this exercise. The model is intentionally explicit so its assumptions are inspectable rather than hidden behind a vendor client.
1. A distributed database is a set of independent participants
A process is a running program with its own memory and execution state. A node is the logical database participant that runs one or more processes and is usually associated with a host/virtual machine/container identity plus persistent storage and network endpoints. Product documentation sometimes uses node, server, member, instance, or replica-set member differently, so always check the product's terminology. In this course, a node means an independently addressed database participant.
A cluster is the cooperating set of nodes that collectively provides a service. The important word is cooperating: cluster membership does not create shared memory or instantaneous knowledge. Nodes exchange messages over a network, may hold different local state at a given instant, and can fail independently or together.
| Term | Course mental model | What it is not |
|---|---|---|
| Process | A running database/service program with memory, sockets, threads/tasks, and failure state | A durable copy of the data by itself |
| Node | An independently addressed database participant with process + local persistent state | Automatically an independent failure domain |
| Cluster | Nodes cooperating to expose a logical service | A single machine with magically shared state |
| Partition / shard | A subset of the logical keyspace or dataset assigned as a unit of placement/routing | A replica/copy |
| Replica | Another copy of a partition or logical data item | A different partition |
| Coordinator | The participant that receives/orchestrates one request in a given protocol | Necessarily the permanent leader |
| Leader / follower | Roles used by leader-based protocols for ordering/replication | Universal roles present in every distributed database |
| Client / proxy | Request originator or routing intermediary | A source of omniscient, always-current cluster truth |
2. Partitioning and replication answer different questions
Partitioning answers “which subset of the
dataset is responsible for this key or range?” If AtlasMart has
four logical partitions p0 through p3,
one order key maps to one of them according to a routing
function or directory. Partitioning can increase capacity and
parallelism because different keys can be served by different
machines, but it creates routing and cross-partition-operation
costs.
Replication answers “which additional nodes store copies of that same partition or value?” Three replicas can improve availability, read capacity, recovery options, or durability, but only under the actual acknowledgement, synchronization, placement, and failure assumptions. Three copies on three processes that share one host or disk controller are not equivalent to three independent failure domains; Lesson 4 makes that concrete.
A replication factor is simply the number of intended copies in a particular configuration. It does not, by itself, state when a write is acknowledged, which replica serves a read, how conflicts are handled, whether replicas are synchronously current, or what failures can be tolerated. Those are separate protocol and topology choices.
3. A coordinator is a request role, not necessarily a permanent owner
Many systems accept a request at a node and have that node coordinate the operation. The coordinator may look up routing metadata, forward the operation to the right owners, collect acknowledgements, and return a result. In some leader-based designs, a client or proxy must reach the current leader for a partition. In leaderless/Dynamo-style designs, a coordinator can be a temporary role for that request. In proxy-based systems, a stateless router may perform part of this work. Therefore “the coordinator” is not a universal synonym for “leader.”
Routing metadata records which node(s) currently own a partition, token range, shard, tablet, or similar unit. Because topology changes, routing metadata itself has a version. This chapter calls that version an epoch: a monotonically newer generation number used to detect stale ownership information. Different products may use terms, generations, configuration versions, leases, or epochs with different semantics; the mechanism is the important part.
If a client cached “p3 belongs to n4” yesterday, but p3 moved today, blindly sending writes according to the old map can create errors or, in a badly designed system, divergent ownership. A safe implementation needs a way to detect stale metadata and redirect/reload rather than assuming topology is static.
4. Trace one AtlasMart write end to end
The lab uses the key tenant-normal-042:order-1001.
A deterministic CRC32 function maps it to one of four
partitions. Partition p3 is represented by replicas
n4, n5, and n6;
n4 acts as coordinator in this simplified topology.
The demonstration sets an acknowledgement threshold of two
durable local appends. That number is a
lab policy only; it is not presented as a
universal quorum formula or a proof of linearizability.
Each replica has a simplified
write-ahead log (WAL): durable log evidence
written before later in-memory/index/storage reorganization.
Real storage engines differ, but the lab needs one observable
point at which a replica can claim it durably recorded version
7. Node n6 is temporarily delayed and remains at
version 6 when the client receives its acknowledgement. A
background repair step later sends the missing version.
from dataclasses import dataclassfrom zlib import crc32@dataclass(frozen=True)class Route: partition: str coordinator: str replicas: tuple[str, ...]EPOCH = 17ROUTES = { "p0": Route("p0", "n1", ("n1", "n2", "n3")), "p1": Route("p1", "n2", ("n2", "n3", "n4")), "p2": Route("p2", "n3", ("n3", "n4", "n5")), "p3": Route("p3", "n4", ("n4", "n5", "n6")),}def partition_for(key: str) -> str: return f"p{crc32(key.encode()) % len(ROUTES)}"key = "tenant-normal-042:order-1001"partition = partition_for(key)route = ROUTES[partition]version = 7old_version = 6ack_threshold = 2replica_versions = {node: old_version for node in route.replicas}wal = {node: [] for node in route.replicas}timeline = []timeline.append((0, "client-east", f"lookup route epoch={EPOCH}: {key} -> {partition}"))timeline.append((2, "client-east", f"PUT v{version} -> coordinator {route.coordinator}"))wal[route.coordinator].append(version)replica_versions[route.coordinator] = versiontimeline.append((4, route.coordinator, f"append WAL v{version}; forward to {route.replicas[1:]}"))# n4 acknowledges promptly; n5 is temporarily unreachable from the coordinator.fast_replica = route.replicas[1]slow_replica = route.replicas[2]wal[fast_replica].append(version)replica_versions[fast_replica] = versiontimeline.append((10, fast_replica, f"append WAL v{version}; ACK"))timeline.append((12, slow_replica, "message delayed; remains at v6"))ack_set = {route.coordinator, fast_replica}if len(ack_set) >= ack_threshold: timeline.append((13, route.coordinator, f"client ACK after durable ACKs={sorted(ack_set)}"))before_repair = dict(replica_versions)wal[slow_replica].append(version)replica_versions[slow_replica] = versiontimeline.append((200, "repair-worker", f"compare {partition}; send missing v{version} to {slow_replica}"))timeline.append((205, slow_replica, f"append WAL v{version}; converge"))print(f"routing_epoch={EPOCH}")print(f"key={key}")print(f"partition={partition} coordinator={route.coordinator} replicas={','.join(route.replicas)}")print(f"ack_threshold={ack_threshold}")print("timeline:")for ms, actor, event in timeline: print(f" t+{ms:03d}ms {actor:13s} {event}")print(f"ack_set={sorted(ack_set)}")print(f"before_repair={before_repair}")print(f"after_repair={replica_versions}")print("note=ack threshold is a demo policy; by itself it does not prove linearizability or protection from correlated failure")
Verified deterministic output
routing_epoch=17key=tenant-normal-042:order-1001partition=p3 coordinator=n4 replicas=n4,n5,n6ack_threshold=2timeline: t+000ms client-east lookup route epoch=17: tenant-normal-042:order-1001 -> p3 t+002ms client-east PUT v7 -> coordinator n4 t+004ms n4 append WAL v7; forward to ('n5', 'n6') t+010ms n5 append WAL v7; ACK t+012ms n6 message delayed; remains at v6 t+013ms n4 client ACK after durable ACKs=['n4', 'n5'] t+200ms repair-worker compare p3; send missing v7 to n6 t+205ms n6 append WAL v7; convergeack_set=['n4', 'n5']before_repair={'n4': 7, 'n5': 7, 'n6': 6}after_repair={'n4': 7, 'n5': 7, 'n6': 7}note=ack threshold is a demo policy; by itself it does not prove linearizability or protection from correlated failure
The important evidence is not that “replication worked.” It is
the exact state at each observation point. At client
acknowledgement time, n4 and n5 hold
version 7 while n6 is still at version 6. After
repair, all three converge. That distinction matters later when
we discuss stale reads, consistency levels, repair, failure
during acknowledgement, and partitions.
5. What this evidence proves—and what it does not
The timeline proves only the rules encoded by the simulator: the client routed using epoch 17; one coordinator and one additional replica durably appended version 7 before acknowledgement; a third copy lagged; background repair later copied version 7. It does not prove that reads cannot return version 6, that two copies survive a rack failure, that a concurrent writer cannot create a conflicting version, or that an acknowledged write is linearizable.
This is a general distributed-systems discipline: separate observed events from protocol guarantees. “I saw two ACKs” is evidence. “Therefore the system is strongly consistent” is an inference that requires additional assumptions about ownership, ordering, failure, quorum intersection, read rules, epochs, and conflict handling. Those assumptions are the subject of later chapters.
6. Deliberately wrong approach: bind one logical key to one hard-coded server
A common first implementation puts orders on
n4 and stores that endpoint in application
configuration. It appears simple until n4 is
drained, partition ownership moves, or the client runs in a
network location that cannot reach n4. The
existence of replicas does not help if the client has no
discovery/routing mechanism and no way to refresh stale
metadata.
The safer pattern is explicit indirection: use a product-supported driver/router/proxy, versioned routing directory, or protocol-aware redirection. Cache routing metadata for efficiency only when the client can detect that it is stale. Log the partition/shard and target node for diagnostics, but avoid making internal topology a permanent application data model.
Routing metadata can reveal hostnames, addresses, tenant placement, and cluster structure. Treat administrative topology endpoints as privileged information. Application clients should receive only the routing capability they require; “knowing the shard” must not become authorization to read another tenant's data.
7. Production judgment and bridge to partial failure
Use the topology vocabulary to ask concrete questions: Which component owns routing? How quickly can a stale route be detected? Is a coordinator a permanent role or per-request role? Which local event counts as durable? Which replicas may lag? How are lag and repair exposed? Which data paths are single-partition and which become scatter/gather or cross-partition operations? How is tenant placement protected?
The next lesson removes another single-machine assumption: there is no single “up/down” bit for a distributed system. One node can be slow, another can have failed storage, two observers can have asymmetric network reachability, and an unrelated partition can continue working. We will make those contradictory-looking states visible.
Verification checklist
- The client key maps deterministically to one partition.
- The partition map contains an explicit routing epoch.
- The output distinguishes partition from replica and coordinator from leader.
- The client acknowledgement occurs while one replica is still stale.
- Background repair visibly changes the lagging replica from version 6 to version 7.
- The lesson does not claim that two acknowledgements automatically imply linearizability or correlated-failure tolerance.
Check your understanding
- What is the difference between a partition and a replica?
- Why can a coordinator be different from a leader?
- What problem does an epoch or configuration version solve?
- At the instant the lab acknowledges version 7, what does n6 contain?
- Why does observing two acknowledgements not prove a universal consistency guarantee?
Review the answers
A partition is a unit/subset of data placement; a replica is an additional copy of the same partition or logical data.
Coordinator is a role in handling one request. Some protocols route to a leader, while others allow different nodes/proxies to coordinate requests.
It lets participants detect that cached ownership/routing information belongs to an older topology and must be refreshed or redirected.
n6 still contains version 6 in the simplified model; repair later advances it to version 7.
The guarantee also depends on read rules, ordering, ownership, failure domains, concurrent writes, acknowledgement semantics, and protocol assumptions that the count alone does not state.
Authoritative references
- A Note on Distributed Computing — Sun Microsystems Laboratories paper emphasizing latency, concurrency, and partial failure as fundamental differences from local calls.
- Dynamo: Amazon's Highly Available Key-value Store — Primary paper describing partitioning, coordinators, replication, versioning, and availability-oriented request paths.
- Bigtable: A Distributed Storage System for Structured Data — Primary paper illustrating partition/tablet placement and distributed storage responsibilities.
- Python 3.14.7 release — Official Python release used to re-check the current stable standard-library baseline at chapter-generation time.