Chapter 08 · Membership, Failure Detection, Gossip, and Cluster Coordination
Static Membership vs Dynamic Discovery and Service Registration
Build a failure-aware membership model that separates discovery endpoints from persistent node identity and data ownership, then prove stale discovery can be rejected with topology epochs.
Learning outcomes
AtlasMart is replacing node-b during a
catalog-shard expansion. The old process is gone, a new process
starts on the same address, one client still holds stale
discovery data, and the routing table is changing at the same
time. This is the point where “a cluster is a list of IP
addresses” stops being a useful model.
Distinguish a process, persistent node identity, endpoint, seed, discovery record, membership view, and data-ownership map.
Explain bootstrap, join, leave, replacement, and why discovery alone must not silently transfer ownership.
Trace one AtlasMart request from client discovery through routing metadata to the current shard owner.
Diagnose stale service-discovery data and address reuse without confusing reachability with identity.
Design safe membership changes around identity, topology epochs, ownership convergence, and rollback evidence.
Chapter 07 treated partition maps and routing epochs as first-class state. Membership is the machinery that tells clients and nodes which processes exist; ownership metadata separately tells them which data those processes are authorized to serve. Conflating those two maps is a split-brain precursor.
1. A member is more than an address
A process is a running program instance. A node is the durable cluster identity the system assigns to a database participant; many systems persist an identifier so a restart is distinguishable from an unrelated process. An endpoint is a network address used to contact that participant. A cluster is the cooperating set of nodes plus the protocols and metadata that define membership and ownership.
A seed is a bootstrap contact, not necessarily a leader and not necessarily the authoritative owner of membership. Service discovery maps a service name or registration query to endpoints. A membership view is a node's current belief about members and their states. A routing/ownership map maps keys or partitions to the nodes allowed to coordinate or store them. These maps often converge through different mechanisms and on different schedules.
| Object | Example | Can change independently? |
|---|---|---|
| Persistent node identity | node-b |
Should survive ordinary restart; replacement normally gets a new identity or explicit replace operation |
| Endpoint | 10.0.0.12:9042 |
Yes; addresses can move or be reused |
| Discovery record | service → endpoint set | Yes; DNS/registries cache and expire |
| Membership view | ALIVE / LEAVING / LEFT / SUSPECT | Yes; views can temporarily disagree |
| Ownership map | partition 17 → node-b | Yes; rebalancing changes ownership under a topology epoch |
2. Static seeds, dynamic discovery, and registration solve different problems
A fixed seed list is operationally simple but becomes stale when addresses change. DNS or cloud discovery reduces manual endpoint management, but its records have cache/TTL and propagation semantics. A service registry can add health and metadata, yet its own availability and consistency now matter. Peer-to-peer membership can discover changes without a central lookup on every request, but each node temporarily holds a local view rather than instantaneous global truth.
The safe design question is therefore not “which discovery technology is best?” It is “what state is authoritative for identity and ownership, how is it versioned, and what must a client verify before sending a write?”
Discovery answers “where might I connect?” It should not, by itself, answer “who owns partition 17 right now?” A topology epoch, term, generation, or equivalent ownership version lets stale clients detect that their routing metadata has expired.
3. Trace one request through discovery and ownership
client cache: service=atlas-catalog endpoints=[10.0.0.11,10.0.0.12]
client routing epoch: 18
key: catalog:sku-442 -> partition catalog-17
owner at epoch 18: node-b @ 10.0.0.12
maintenance:
node-b -> LEAVING -> LEFT
node-d bootstraps at 10.0.0.12
new ownership epoch: 19, catalog-17 -> node-d
stale client:
connects 10.0.0.12 successfully
BUT endpoint is now node-d, not node-b
safe response: WRONG_EPOCH / refresh route
The successful TCP connection proves reachability only. It does not prove that the process has the identity, generation, shard lease, or data needed by the stale route. A protocol that returns node identity and topology version lets the client distinguish “address reachable” from “route valid.”
4. Wrong approach: treat endpoint reuse as node continuity
A real operational shortcut is to replace a failed VM and deliberately reuse its IP address so old clients “keep working.” Without an explicit replacement protocol, persistent identity, and ownership fencing, the new process may not possess the old node's complete data or authority. A stale client can then interpret any acknowledgement from that endpoint as proof that the old owner survived.
The repair is to make replacement an explicit topology operation: bootstrap the new identity, stream/restore required state, publish the new membership view, advance ownership metadata, wait for routing convergence, and only then retire compatibility routes. If a product supports an identity-preserving replace operation, use its documented preconditions rather than simulating it by DNS/IP reuse.
5. AtlasMart lab — stale discovery versus persistent identity
This simulation requires only Python 3.11+ standard library. It does not perform DNS changes or open sockets. It models a stale registry entry and an address that has been reused by a different node identity.
from dataclasses import dataclass
@dataclass(frozen=True)
class Endpoint:
node_id: str
address: str
topology_epoch: int
# AtlasMart shard catalog at routing epoch 18.
owner = {"partition": "catalog-17", "node_id": "node-b", "epoch": 18}
# node-b was decommissioned. A new process reuses the same address.
registry_fresh = Endpoint("node-d", "10.0.0.12:9042", 19)
registry_stale = Endpoint("node-b", "10.0.0.12:9042", 18)
actual_process = registry_fresh
print("owner-map:", owner)
print("stale discovery:", registry_stale)
print("actual process:", actual_process)
# WRONG: endpoint address is treated as identity.
naive_accepts = registry_stale.address == actual_process.address
print("naive endpoint-only identity accepts:", naive_accepts)
if naive_accepts:
print("FAILURE: response from node-d could be mistaken for retired node-b")
# SAFER: validate persistent node identity and topology epoch.
safe_accepts = (
actual_process.node_id == owner["node_id"]
and actual_process.topology_epoch == owner["epoch"]
)
print("safe identity+epoch accepts:", safe_accepts)
if not safe_accepts:
print("ACTION: reject stale route, refresh membership/ownership metadata")
# After coordinated ownership transfer.
owner = {"partition": "catalog-17", "node_id": "node-d", "epoch": 19}
print("refreshed owner-map:", owner)
print("safe retry accepted:", actual_process.node_id == owner["node_id"] and actual_process.topology_epoch == owner["epoch"])
Expected output
owner-map: {'partition': 'catalog-17', 'node_id': 'node-b', 'epoch': 18}
stale discovery: Endpoint(node_id='node-b', address='10.0.0.12:9042', topology_epoch=18)
actual process: Endpoint(node_id='node-d', address='10.0.0.12:9042', topology_epoch=19)
naive endpoint-only identity accepts: True
FAILURE: response from node-d could be mistaken for retired node-b
safe identity+epoch accepts: False
ACTION: reject stale route, refresh membership/ownership metadata
refreshed owner-map: {'partition': 'catalog-17', 'node_id': 'node-d', 'epoch': 19}
safe retry accepted: True
The first Boolean is intentionally True:
endpoint-only identity accepts the wrong process. The
identity+epoch check rejects it, forces metadata refresh, and
only then accepts the new owner. The lab proves the value of
explicit identity/version checks; it does not model a specific
database's join protocol.
Check your understanding
- Why is a seed node not automatically a leader?
- What does a successful connection to a reused IP prove?
- Why coordinate membership with ownership?
- What is the purpose of a topology epoch?
- When can static seeds be reasonable?
Review the answers
1. A seed is normally just a known bootstrap contact. Leadership/coordination is a separate protocol role.
2. Only that something is reachable at that endpoint; it does not prove persistent node identity or shard ownership.
3. Because adding/removing a process changes where replicas may safely live; clients must not route writes to a node before its state/authority is ready.
4. It gives clients/nodes a monotonic version for detecting stale routing or ownership metadata.
5. Small, stable environments where bootstrap endpoints are intentionally managed and the seed list is not mistaken for the live ownership map.
6. Production judgment
Choose discovery by operational environment, but make identity and ownership explicit. Observe join duration, membership-view divergence, stale discovery records, rejected stale epochs, bootstrap/stream progress, and client route-refresh rates. Protect registries and administrative membership APIs with authentication, least privilege, and audit logging: poisoning discovery or forcing membership changes is a control-plane security event.
Do not change several identities and ownership ranges simultaneously unless the system has enough replica/capacity headroom and a documented rollback path. This leads directly to failure detection: once nodes know who should exist, they still need a disciplined way to decide whether a silent peer is merely slow or probably unavailable.
Authoritative references
- SWIM membership protocol — separates failure detection from membership dissemination and formalizes weakly consistent membership
- Apache Cassandra 5.0.9 download — current implementation-version snapshot checked August 2026; optional reference only
- Cassandra nodetool gossipinfo — example of observable gossip/member state in a current distributed database