Membership Changes and Replica Set Evolution

LESSON

Consistency and Replication

017 30 min advanced

Membership Changes and Replica Set Evolution

By the end of this lesson, you will be able to...

  • distinguish a failed-heartbeat suspicion from an authoritative membership change;

  • trace how a replacement replica becomes safe to vote and serve;

  • identify why fencing and configuration epochs protect a cluster when an old node returns.

Idea in one sentence: A timeout says that communication is missing; only a committed configuration changes who may vote, replicate, or serve.

Core Insight

Shard 184 has three voters: md-db-2, md-db-4, and ny-db-3. During a transatlantic congestion burst, New York stops acknowledging heartbeats. Madrid can observe silence. It cannot observe the cause: a crash, a pause, or delayed packets all look similar for a while.

The tempting reaction is to immediately stop counting New York and add a replacement. That can make a healthy but slow node disappear from one machine's quorum calculation before the cluster has agreed on a new quorum. The result is not faster recovery; it is conflicting authority.

The safe model has four distinct stages:

suspicion -> committed configuration -> catch-up -> fenced participation

Failure detection supplies evidence. Reconfiguration records the cluster's decision. Catch-up gives a replacement the committed history it needs. Fencing keeps old membership from leaking back into the new configuration.

The Initial Model: Timeout Means Removal

A short timeout is attractive because it detects genuine crashes quickly. If ny-db-3 has really stopped, waiting for it can keep a three-voter group from reaching a majority after another failure.

But a timeout does not prove a process is dead. A stalled virtual machine, a congested link, a garbage-collection pause, and a power loss can all produce the same missing heartbeat. The initial model works only when the network and hosts are healthy enough that a timeout is strong evidence. Production systems do not get that condition for free.

The trade-off is explicit. A 500 ms detector reacts sooner but can trigger false suspicion during a delay burst. A 3 s detector reduces churn but leaves a real failure in the quorum picture longer. The detector's job is to choose a policy response—probe, alert, prepare replacement—not to edit quorum math on its own.

The Moving Parts

Keep these objects separate:

Object What it says What it cannot authorize alone
Suspicion A peer has missed expected communication. A new voting quorum.
Term or leader state Which leader and ordering epoch are current. A changed replica set.
Committed membership epoch Which nodes count for quorum now. That a new node has all needed data.
Learner or non-voter A node may receive history and prove catch-up. Serving or voting as a full replica.
Fencing token / epoch check A request belongs to the current configuration. Repair of a stale node by itself.

The term membership means more than a discovery list. It is control-plane state that changes the data plane: who receives quorum appends, which acknowledgements count, which routes are valid, and who may use leader or lease authority.

The Mechanism Step by Step

Assume ny-db-3 remains unreachable and Harbor Point has warm standby sg-db-1. The following trace is illustrative; product-specific protocols differ in their exact transition, but they need one authoritative ordering of membership decisions.

epoch 81 voters = {md-db-2, md-db-4, ny-db-3}

1. The detector marks ny-db-3 suspect after repeated missed probes.
   Epoch 81 is still the quorum rule.

2. The leader adds sg-db-1 as a learner in committed control state.
   It receives log entries or a snapshot, but its acknowledgement does not yet count.

3. sg-db-1 proves it has reached the required committed position.
   The leader can now propose a voter change.

4. The cluster commits a transitional configuration, or an equivalent
   protocol that preserves quorum intersection across the change.

5. The new voter configuration is committed:
   epoch 84 voters = {md-db-2, md-db-4, sg-db-1}.

6. Routers, replication workers, and lease checks accept epoch 84.
   ny-db-3 is fenced until it rejoins through the new workflow.

The intermediate state is the hard part. Promoting sg-db-1 before it catches up increases the number of voters without adding a useful copy. Different nodes may then need acknowledgements from different groups while the replacement still lacks the history that makes its acknowledgement meaningful.

Many consensus systems use a joint or overlapping configuration during the handoff. The teaching rule is simpler: a quorum that could commit before the change must intersect with a quorum that can commit after it. This prevents two disjoint groups from each treating themselves as the continuing cluster.

So far: the detector did not change membership. The committed configuration did. The learner made the replacement useful before its vote mattered. The transition prevented old and new quorum rules from passing each other in the dark.

The Unsafe Shortcut, Made Visible

Suppose md-db-2 takes a shortcut after the first timeout. It locally decides that ny-db-3 is gone and starts counting sg-db-1 immediately. Meanwhile md-db-4 has not received that decision and still uses epoch 81.

md-db-2 believes voters are {md-db-2, md-db-4, sg-db-1}
md-db-4 believes voters are {md-db-2, md-db-4, ny-db-3}
ny-db-3, after a delay, still believes epoch 81 exists

Each view looks plausible from one machine. Together they are dangerous. Acknowledgements no longer have one unambiguous meaning. A client route might send traffic to a node that another replica considers removed. A repair task might copy data according to the wrong owner list. The control plane has become less reliable than the user data it is meant to protect.

The committed transition fixes the reasoning order. First, the cluster records the membership state through its normal authority mechanism. Then every dependent path reads that one result.

Dependent path Question it asks after a committed epoch
Quorum replication Which acknowledgements count for this log entry?
Client routing Which leader and replicas may receive this shard's traffic?
Lease or safe read Which term and configuration back this read proof?
Catch-up worker Is the target a learner, a voter, or a removed node?
Repair process Which replicas are valid peers for this configuration?

This does not mean a configuration change is free from failure. It can be delayed because the current quorum cannot commit it. The learner can fail during a snapshot. A removed node can return through a stale load balancer. The mechanism does not eliminate these cases; it makes each case have a named state and a safe next action.

Operational Failure Modes

A WAN jitter event repeatedly removes and re-adds the same replica. The detector threshold is treating temporary delay as permission to reconfigure. Add a suspicion grace period or additional probes, and require stability before starting the reverse transition. Fast automation is useful only when it does not create more recovery work than the interruption did.

Adding a replacement causes leader CPU and write latency to spike. The leader is sending a large state transfer on the same path that handles quorum appends. Keep the replacement non-voting, apply the bounded recovery policy from the flow-control lesson, and promote only after catch-up evidence is available.

A removed node still serves a stale read. Some route, proxy, or RPC path is not checking the configuration generation. Treat stale-epoch rejects as a first-class signal and quarantine the node until it has rejoined through the current configuration.

When the Old Node Returns

At 10:15, ny-db-3 returns. It still has files from epoch 81 and an old route may still point at it. A hostname and a healthy process are not proof of current authority.

Fencing makes the difference visible. Requests, replication messages, client routes, and lease decisions include the current term or configuration generation. A node that presents an obsolete epoch rejects service or is rejected by the live cluster. It cannot become a shadow voter just because it remembers an earlier membership list.

Only after the returning node learns the current configuration and catches up from the committed log or a snapshot should it be eligible to participate again. This protects read correctness too: an old node must not serve a lease-based or “current” read after the cluster has advanced to a new configuration.

There is a useful diagnostic distinction here. A catch-up failure means the node is not yet useful as a replica. A fencing failure means the node can still affect clients or peers while it is not useful. The first consumes recovery capacity; the second threatens authority. Keep separate metrics for catch-up position, stale-epoch request rejects, and traffic served by retired routes so operators do not mistake one problem for the other.

Cost, Limits, and Signals

This workflow improves safety, but it is not free. Catch-up spends leader network, disk, and snapshot capacity. Conservative removal leaves a suspected voter in the group longer. Aggressive automation reduces manual work but can flap during intermittent WAN failures and repeatedly trigger expensive transfers.

Watch signals that select an action rather than merely describe a node:

Signal Decision it supports
Consecutive missed probes and probe latency Increase suspicion or begin replacement review.
Current committed membership epoch Decide whose acknowledgement counts.
Learner match or applied position Decide whether promotion is safe.
Snapshot progress and retained-log pressure Choose streaming catch-up or snapshot recovery.
Stale-epoch RPC rejects Detect old routes or unfenced returning nodes.

Membership machinery does not repair an application-level invariant, and it does not make a new replica durable before it catches up. It preserves one authoritative topology so the replication and read guarantees from the rest of this track remain meaningful.

Before an automated change, test the same trace with a delayed but healthy node, a genuinely failed node, and a returning stale node. A controller that handles only the clean replacement path has not yet earned authority over production membership during incidents.

Check Your Understanding

Check: ny-db-3 misses four heartbeats, but no membership entry has been committed. Which voters count now?

Think first, then reveal.

Answer: The voters in the last committed configuration still count: md-db-2, md-db-4, and ny-db-3. The timeout can trigger probes or a proposal, but it is not the topology change.

Check: Why not count sg-db-1 as a voter immediately after it connects?

Think first, then reveal.

Answer: It may not contain the committed history needed for its acknowledgement to be useful. Keep it as a learner while it catches up, then promote through the committed membership protocol.

Practice: Review a Replacement Plan

A five-voter group has one unreachable node and an empty replacement. Write a short plan with one detector action, one committed membership action, one catch-up condition, and one fencing rule.

A good answer should mention: suspicion as evidence rather than removal; a learner or equivalent pre-voter role; a required committed position or snapshot validation before promotion; and epoch-aware routes and RPCs that reject the old node until it rejoins.

Connections

Resources

Key Takeaways

  1. A timeout creates suspicion, not a new quorum.
  2. Safe replacement separates catch-up from voting and commits the topology transition before routing changes.
  3. Fencing by term and configuration epoch keeps a returning node's old state from becoming old authority.
PREVIOUS Clocks, Leases, and Safe Reads NEXT Replication Flow Control and Backpressure