Chapter 08 · Membership, Failure Detection, Gossip, and Cluster Coordination
Operational Coordination: Topology Changes, Rolling Maintenance, and Convergence Windows
Plan topology changes as checkpointed state transitions that preserve replica/quorum redundancy through rolling maintenance and convergence windows.
Learning outcomes
AtlasMart must replace one database node per failure domain without losing the write quorum for any RF=3 partition. The unsafe runbook says “take the old servers down and let the cluster heal.” The safe runbook advances only after membership, routing, replica catch-up, and repair evidence converge.
Plan rolling maintenance around failure domains and replica placement rather than node count alone.
Define convergence checkpoints for membership, routing epochs, data streaming/repair, and client health.
Explain why simultaneous topology changes can erase redundancy even when the cluster still has many live machines.
Separate drain/decommission, replacement/bootstrap, catch-up, validation, and next-node authorization.
Write a rollback-aware maintenance sequence with explicit SLO/error-budget and security considerations.
Resharding already required routing-version changes, copy/catch-up, validation, and headroom. Rolling membership changes add failure-detector/gossip convergence and the risk that operators create a second failure before the first change has finished healing.
1. A rolling operation is a state machine
“Rolling” means changes are staged so the system preserves required service and durability while one unit is unavailable or transitioning. It does not mean “restart machines sequentially as fast as possible.” A safe step has entry criteria, observable transition state, exit criteria, and rollback.
| Phase | Evidence before advancing | Typical rollback question |
|---|---|---|
| preflight | healthy replica placement; error budget; backups; no active repair emergency | can we cancel before reducing redundancy? |
| drain/fence old node | client traffic removed; ownership generation protected | can old node safely resume if replacement fails? |
| join/bootstrap replacement | persistent identity accepted; stream/catch-up running | can replacement be removed cleanly? |
| converge | membership/routing views agree enough; repair backlog acceptable/zero | which epoch is authoritative? |
| validate | quorum/read/write SLOs, replica counts, logs, checksums/history | can traffic remain on old topology? |
| next failure domain | all prior exit criteria satisfied | if another node fails now, do we still meet quorum? |
2. Failure-domain math matters more than total node count
Suppose each AtlasMart partition has replication factor (RF) 3 with one replica in each of three zones and the write policy needs two acknowledgements. Taking one zone's replica down leaves two available copies. Taking a second zone's copy before the first has been replaced/caught up leaves only one and makes writes unavailable for that partition.
A six-node cluster may still be unsafe if the relevant partition's three replicas are concentrated on the exact nodes being maintained. Therefore the runbook should query/verify per-partition placement and failure-domain health, not only “5 of 6 nodes are up.”
3. Define convergence checkpoints
Before touching next node/failure domain:
[ ] membership view contains expected identities/states
[ ] routing/topology epoch has advanced and stale routes are declining
[ ] replacement owns intended ranges only after bootstrap/catch-up
[ ] replica count/failure-domain placement meets policy
[ ] repair/stream backlog is zero or below an explicitly safe threshold
[ ] write/read quorum succeeds from representative clients
[ ] p95/p99 latency and error rate are inside maintenance SLO
[ ] stale-owner/fencing rejections are understood
[ ] backups/rollback artifacts remain valid
[ ] no unrelated incident is consuming the remaining redundancy
These checks are deliberately mechanism-oriented. “Node is green in dashboard” is not enough if its data is incomplete, routing metadata is stale elsewhere, or the remaining replicas share one failure domain.
4. Wrong approach: overlap topology changes to save time
Operators often accelerate a maintenance window by starting node B while node A is still streaming/repairing. The visible cluster may look healthy because both replacements answer health checks, yet some partitions still have only one fully current replica. If a real fault occurs during this window, the cluster can lose quorum or promote stale state.
The safe correction is to serialize by the redundancy budget, not necessarily by host: change only as many members/failure domains as the placement and quorum policy can tolerate, throttle data movement to protect foreground SLOs, and wait for explicit convergence evidence before consuming more redundancy.
5. AtlasMart lab — unsafe overlap versus checkpointed rolling change
The simulation models one RF=3 partition with conceptual write
quorum W=2. The unsafe sequence takes a second zone
offline before restoring the first. The safe sequence
restores/catches up each zone before advancing.
from dataclasses import dataclass
@dataclass
class Step:
name: str
available_replicas: int
routing_epoch: int
repair_backlog: int
W = 2 # conceptual write quorum for an RF=3 partition
unsafe = [
Step("start", 3, 50, 0),
Step("zone-a node offline", 2, 50, 12),
Step("operator also takes zone-b node offline", 1, 50, 12),
]
print("UNSAFE rolling sequence")
for s in unsafe:
print(s.name, "available=", s.available_replicas, "write_quorum=", s.available_replicas >= W)
safe = [
Step("start", 3, 50, 0),
Step("drain a1", 2, 50, 8),
Step("a2 joins and streams", 2, 51, 3),
Step("a2 caught up; routing converged", 3, 52, 0),
Step("only now drain b1", 2, 52, 7),
Step("b2 caught up; routing converged", 3, 53, 0),
]
print("\nSAFE checkpointed sequence")
for s in safe:
checkpoint = s.available_replicas >= W and (s.repair_backlog == 0 or s.name.startswith("drain") or "streams" in s.name)
print(f"{s.name:35} available={s.available_replicas} epoch={s.routing_epoch} repair={s.repair_backlog} quorum={s.available_replicas >= W}")
print("\nRule: do not advance to the next failure domain until membership/routing and replica catch-up evidence converge.")
Expected output
UNSAFE rolling sequence
start available= 3 write_quorum= True
zone-a node offline available= 2 write_quorum= True
operator also takes zone-b node offline available= 1 write_quorum= False
SAFE checkpointed sequence
start available=3 epoch=50 repair=0 quorum=True
drain a1 available=2 epoch=50 repair=8 quorum=True
a2 joins and streams available=2 epoch=51 repair=3 quorum=True
a2 caught up; routing converged available=3 epoch=52 repair=0 quorum=True
only now drain b1 available=2 epoch=52 repair=7 quorum=True
b2 caught up; routing converged available=3 epoch=53 repair=0 quorum=True
Rule: do not advance to the next failure domain until membership/routing and replica catch-up evidence converge.
The failure signal is explicit: the unsafe sequence reaches
write_quorum=False. The safe sequence keeps at
least two replicas available and records routing epochs plus
repair backlog so “process alive” is not confused with “replica
ready.”
Check your understanding
- Why is node count an insufficient maintenance metric?
- What is a convergence window?
- Why wait for repair/stream evidence?
- When can parallel maintenance be safe?
- What should trigger rollback?
Review the answers
1. Quorum and durability depend on the replica set and failure-domain placement for each partition, not aggregate live hosts.
2. The period after a topology change when membership, routing, data placement, and client views have not all caught up.
3. A replacement can be reachable before it contains the required replica state.
4. Only when placement/quorum analysis proves the simultaneous changes do not consume required redundancy and the system has headroom.
5. Explicit SLO, quorum, convergence, data-integrity, or security thresholds—not operator intuition alone.
6. Production judgment and chapter bridge
A production runbook should identify owner, maintenance window, version/edition compatibility, backup/restore point, failure domains, per-partition replica policy, client retry behavior, expected membership/routing transitions, throttles, abort thresholds, and rollback. Use canary nodes when meaningful, but do not infer cluster-wide safety from one canary if partition placement differs.
Protect maintenance credentials and administrative endpoints, record who forced membership changes, and avoid disabling gossip/failure detection as a casual workaround. Current Cassandra documentation, for example, exposes explicit bootstrap, gossip, and failure-detector operational tools; those are product-specific controls with real blast radius, not generic lab commands.
With membership, suspicion, dissemination, fencing, and rolling coordination established, Chapter 09 moves below the replication protocol into storage-engine mechanics: write-ahead logs, B-trees, LSM trees, SSTables, compaction, and amplification.
Authoritative references
- Apache Cassandra 5.0.9 download — current GA snapshot checked August 2026
- Cassandra nodetool reference — product-specific bootstrap/removal/repair observability and controls
- SWIM membership protocol — membership convergence and suspicion model
- Raft paper — term/quorum coordination background for ownership changes