skip to content

A Redis Cluster primary that owns a range of hash slots has been marked as failed. Describe how one of its replicas becomes the new owner of those slots.

level: seniorimportance: must knowfreq 48%

answer

  1. FAIL first — never self-promote on a lost link
  2. Eligibility: link-down < node-timeout × validity-factor
  3. Delay = 500 ms + random + rank×1 s, rank by offset
  4. Majority of slot-owning primaries vote once per epoch
  5. Highest configEpoch wins; old primary returns as replica

basics

~20 s

Eligible replicas wait a short delay ranked by how far behind they are, then request votes for a new configuration epoch. Primaries each vote once per epoch; a replica winning a majority bumps its config epoch, claims the slots and announces itself, and the old primary rejoins as its replica.

solid answer

~1 min

Once the primary is in the FAIL state, each of its replicas checks **eligibility**: it must have been connected to the primary recently enough (a link down for longer than roughly `cluster-node-timeout × cluster-replica-validity-factor` disqualifies it) so a hopelessly stale replica cannot win. Eligible replicas then wait before campaigning. The delay is a small fixed part, a random part to avoid collisions, and a **rank** component derived from replication offset: the replica with the most data has rank 0 and waits least, each further-behind replica adds roughly a second. This biases the outcome toward the freshest replica without needing an explicit comparison protocol. When its delay elapses, the replica increments the cluster's `currentEpoch` and broadcasts a failover authorisation request. Primaries serving slots vote at most once per epoch, and only if the requester's primary really is FAIL. On collecting votes from a **majority of primaries**, the replica takes a new, higher `configEpoch`, claims the failed primary's slots, converts itself to a primary and advertises the new configuration in its pings. Higher config epoch wins conflicts, so the old primary, on returning, reconfigures itself as a replica of the new one. `CLUSTER FAILOVER` performs the same handover manually and, when coordinated, without data loss.

code

text · 14 lines
text
# run ON THE REPLICA that should take over
> CLUSTER FAILOVER
OK      # primary pauses clients, replica syncs to the exact offset,
        # then promotes: no data loss

# primary unreachable: skip coordination, still run the election
> CLUSTER FAILOVER FORCE

# last resort: skip the election entirely (no majority consulted)
> CLUSTER FAILOVER TAKEOVER   # may lose data / diverge history

# eligibility and self-healing knobs
CONFIG GET cluster-replica-validity-factor   # default 10 (0 disables the check)
CONFIG GET cluster-migration-barrier         # default 1

go deeper

for a junior

Say that the replica of the failed primary is promoted, that other primaries have to approve it by majority vote, and that the freshest replica is preferred.

for a middle

Add the ranked delay based on replication offset, the eligibility check for stale replicas, and that the epoch number decides configuration conflicts.

for a senior

Give the full sequence including epochs, per-epoch single voting, slot claiming, the old primary demoting itself, and the manual CLUSTER FAILOVER variants and when each is acceptable.

for a principal

Discuss the tuning surface — validity factor as an availability/consistency dial, migration barrier and spare replicas for post-failover coverage, and the client-side topology-refresh behaviour that decides the user-visible outage.

## Precondition Nothing starts until the failed primary is in the **FAIL** state — the cluster-wide agreement reached when a majority of slot-owning primaries have reported it as unreachable. A replica never promotes itself merely because it lost its own link to the primary; that would be exactly how a network problem turns into two primaries for one slot range. ## Step 1: eligibility A replica disqualifies itself if its data is too old to be trustworthy. The rule is based on how long its replication link to the primary had been broken before the failure: if that disconnection exceeds roughly `cluster-node-timeout × cluster-replica-validity-factor` (default factor 10) plus the ping period, the replica will not campaign. Setting the factor to 0 disables the check, which maximises availability — some replica will always try — at the cost of possibly promoting a badly stale one. This is a real availability-versus-consistency dial. ## Step 2: the ranked delay Eligible replicas do not campaign immediately. Each computes a delay of roughly: ``` 500 ms (fixed) + random 0-500 ms + rank * 1000 ms ``` The fixed part gives the FAIL state time to propagate cluster-wide so voters agree the primary is gone. The random part decorrelates simultaneous campaigns. **Rank** is the replica's position when the primary's replicas are sorted by replication offset, best first: the most up-to-date replica has rank 0 and therefore campaigns first, and it usually wins before the others start. Rank is recomputed as offsets are learned, so a replica that catches up improves its position. This is how Redis biases toward minimal data loss without an explicit "who has the most data" negotiation. ## Step 3: requesting votes The campaigning replica increments the cluster's `currentEpoch` — a monotonically increasing logical clock shared by the cluster — and broadcasts a failover authorisation request stamped with that epoch. A primary grants a vote only if all of the following hold: it serves at least one hash slot; it has not already voted in this epoch; the requester is a replica of a primary currently in the FAIL state; and the request's epoch is not stale. Each voter replies with an acknowledgement. Voters also refuse to vote again for the same failed primary's replicas for a short cooldown, which prevents rapid re-elections from thrashing. Because votes come only from slot-owning primaries and a **majority** is required, a minority partition can never authorise a promotion — the same quorum property that governs FAIL detection. ## Step 4: taking over On winning, the replica: 1. Adopts a new `configEpoch` — the highest so far — which is the tiebreaker for configuration conflicts. 2. Turns itself into a primary and claims all hash slots previously served by the failed node. 3. Advertises the new configuration in its PING/PONG payloads, so the rest of the cluster converges. Configuration conflicts are resolved by **highest config epoch wins**: when the old primary comes back and claims the same slots with an older epoch, every node — including itself — recognises the newer claim, and the old primary demotes itself to a replica of the node that replaced it. Its remaining sibling replicas also re-point to the new primary. Thanks to replication ID handoff (PSYNC2), those replicas can often continue with a partial resynchronisation instead of a full one, avoiding a fork-and-transfer storm right after a failover. ## Manual failover `CLUSTER FAILOVER`, run on a replica, does the same handover deliberately and **without data loss**: the primary pauses clients, the replica waits until it has caught up to the primary's exact offset, and only then does the promotion proceed. Two overrides exist for emergencies: `FORCE` skips the coordination with a primary that is unreachable but still runs the normal election, and `TAKEOVER` skips the election entirely and unilaterally seizes the slots with a bumped epoch. `TAKEOVER` can produce data loss and a divergent history; it is a last resort, not an operational routine. ## Replica migration After a failover the promoted primary may be left with no replica, while some other primary has two. Redis handles this with **replica migration**: a primary with spare replicas donates one to an orphaned primary, subject to `cluster-migration-barrier` (default 1 — never donate if it would leave the donor with fewer than one replica). This is why an odd extra replica in the cluster is genuinely useful: it self-heals coverage after a promotion. ## Timing summary Detection (up to `cluster-node-timeout`) + FAIL agreement propagation + the ranked delay (roughly 0.5-2 s for the best replica) + slot claim + client topology refresh. Only the failed shard's slots are affected; the rest of the cluster keeps serving throughout.

  • How does Redis Cluster avoid promoting a badly out-of-date replica?
    Two mechanisms. First, eligibility: a replica whose replication link had been broken for longer than roughly cluster-node-timeout multiplied by cluster-replica-validity-factor refuses to campaign at all. Second, the ranked delay: replicas are ordered by replication offset, and the one with the most data waits the shortest time, so it normally wins the election before staler siblings even start. Neither is a hard guarantee of zero loss, because replication is asynchronous.
  • What happens when the original primary comes back online after the failover?
    It advertises its old claim on the slots, but its configEpoch is lower than the promoted node's, and in Redis Cluster the highest config epoch wins configuration conflicts. Seeing a newer claim for its slots, the returning node reconfigures itself as a replica of the new primary and synchronises from it. Any writes it accepted that never reached the promoted replica are discarded at that point.

saying these in an interview costs you the question

  • Claiming a replica promotes itself as soon as it loses contact with its primary
  • Thinking replicas vote in the election — only slot-owning primaries do
  • Believing the election guarantees no data loss because the most advanced replica wins
  • Using CLUSTER FAILOVER TAKEOVER as a routine operation without understanding it bypasses the majority
  • Assuming the old primary rejoins as a primary and keeps its writes

context