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.

Beginner → Advanced95–115 minutesrolling-maintenance quorum simulationVendor-neutral · Python 3.13.5 simulatorFree/local · no database or cloud requiredLast reviewed: August 2026

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.

01

Plan rolling maintenance around failure domains and replica placement rather than node count alone.

02

Define convergence checkpoints for membership, routing epochs, data streaming/repair, and client health.

03

Explain why simultaneous topology changes can erase redundancy even when the cluster still has many live machines.

04

Separate drain/decommission, replacement/bootstrap, catch-up, validation, and next-node authorization.

05

Write a rollback-aware maintenance sequence with explicit SLO/error-budget and security considerations.

Continuation from Chapter 07

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

text · safe checkpoint checklist
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.

python · rolling_maintenance.py
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

text · deterministic result
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

  1. Why is node count an insufficient maintenance metric?
  2. What is a convergence window?
  3. Why wait for repair/stream evidence?
  4. When can parallel maintenance be safe?
  5. 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

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.