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.

Beginner100–120 minutesRouting + replica-state timeline labVendor-neutral · Python stdlibLast reviewed: August 2026

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.

01

Define process, node, cluster, shard/partition, replica, coordinator, leader, follower, client, proxy, and routing metadata without treating them as interchangeable words.

02

Trace one AtlasMart key from client routing through a coordinator, replica acknowledgements, persisted write-ahead-log evidence, and later repair.

03

Distinguish partitioning, which assigns different data ranges/keys, from replication, which creates additional copies of the same logical data.

04

Explain why routing metadata needs a version/epoch and why clients cannot safely hard-code one server forever.

05

Interpret acknowledgement and replica state as evidence while avoiding claims that an acknowledgement count alone proves a particular consistency or durability guarantee.

Prerequisite connection

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.

Lab baseline reviewed 29 August 2026

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.

Why version routing state?

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.

python · route one key, collect acknowledgements, then repair a lagging replica
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

text · request timeline and per-replica version evidence
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.

Security boundary

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

  1. What is the difference between a partition and a replica?
  2. Why can a coordinator be different from a leader?
  3. What problem does an epoch or configuration version solve?
  4. At the instant the lab acknowledges version 7, what does n6 contain?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.