skip to content

Cluster Failover

You will learn how the cluster detects a dead master through gossip and PFAIL/FAIL states, elects a replica, and promotes it — plus the async-replication window where acknowledged writes vanish. Interviewers ask about split-brain and lost writes to test honest reasoning about Redis's consistency limits.

part ofRedisoverview, primer and where to startread it →
on this pageshow

questions

5

In Redis Cluster, how do the nodes decide that a primary is actually down rather than just slow, and what does the `cluster-node-timeout` setting control in that process?

level: middleimportance: must knowfreq 55%

answer

  1. Cluster bus = data port + 10000, gossip in PING/PONG
  2. PFAIL = local opinion after cluster-node-timeout
  3. FAIL = majority of slot-owning primaries agree, then broadcast
  4. Needs ≥3 primaries for any majority to exist
  5. FAIL on a primary unlocks replica election

basics

~20 s

Nodes ping each other over the cluster bus and gossip what they see. A node unreachable for longer than cluster-node-timeout is flagged PFAIL (possibly failed); once a majority of primaries report PFAIL for it, one marks it FAIL and broadcasts that, which triggers replica election.

solid answer

~60 s

Every node keeps a link to every other node on the **cluster bus** (the data port plus 10000) and pings a rotating subset roughly every second; the ping payload gossips what that node believes about others. Detection is two-stage: 1. **PFAIL — possibly failed, a local opinion.** If a node gets no PONG from another within `cluster-node-timeout` (default 15000 ms), it marks that node PFAIL. This flag alone does nothing; it only says "I cannot reach it". 2. **FAIL — an agreed fact.** PFAIL flags travel in gossip. When a primary sees PFAIL reports for the same node from a **majority of primaries** (itself included) within a bounded window, it promotes the state to FAIL and broadcasts a FAIL message so the whole cluster converges. Only primaries that own hash slots vote in that count, which is why a cluster needs at least three primaries — with two, no majority survives losing one. Marking a primary FAIL is what makes its replicas eligible to start an election. `cluster-node-timeout` therefore sets both the sensitivity of detection and, roughly, the floor on failover time.

code

text · 16 lines
text
> CLUSTER INFO
cluster_state:ok            # 'fail' when slots are uncovered / no majority
cluster_slots_assigned:16384
cluster_slots_ok:16384
cluster_known_nodes:6
cluster_size:3              # number of primaries serving slots

> CLUSTER NODES
<id> 10.0.1.11:6379@16379 master,fail  - 0 ... connected
<id> 10.0.2.12:6379@16379 myself,master - 0 ... connected 5461-10922
<id> 10.0.3.13:6379@16379 slave,pfail  <master-id> ...

# relevant configuration (same value on every node)
CONFIG GET cluster-node-timeout
1) "cluster-node-timeout"
2) "15000"

go deeper

for a junior

Know the two states by name: PFAIL is one node's suspicion after cluster-node-timeout, FAIL is the agreed verdict once most primaries suspect the same node.

for a middle

Explain the gossip mechanism, the majority-of-primaries requirement, and why three primaries is the minimum for automatic failover.

for a senior

Add the timing breakdown a client sees, the minority-partition write cutoff, and the CLUSTER INFO / CLUSTER NODES fields you would read during an incident.

for a principal

Reason about detection as a quorum property of the topology — failure domains, primary count, and the consistency of cluster-node-timeout across the fleet — and about the availability policy set by cluster-require-full-coverage.

## The cluster bus Redis Cluster nodes talk to each other on a second port — the client port plus 10000 — using a compact binary protocol, separate from the port clients use. Each node keeps a TCP link to every other node, and every second it sends PINGs to a few randomly chosen nodes, preferring ones it has not heard from recently. Each PING/PONG carries a **gossip section**: a handful of entries describing other nodes as the sender currently sees them, including their flags. Nothing here is a consensus protocol for data; it is a membership and configuration protocol. ## Stage one: PFAIL When node A sends a PING to node B and receives no PONG for longer than `cluster-node-timeout`, A flags B as **PFAIL** ("possibly failing"). This is deliberately a *local* opinion. A single node may be wrong for many reasons — a one-way network problem, a saturated NIC, its own scheduling stall — so PFAIL grants no authority. Nothing fails over on PFAIL alone. `cluster-node-timeout` (default 15000 ms) is the only knob that sets this threshold, and it is a cluster-wide concept: every node should be configured with the same value or detection becomes inconsistent. ## Stage two: FAIL PFAIL flags propagate through gossip. A primary collects the reports it has heard about a given node and, when the number of distinct primaries reporting PFAIL (counting itself) exceeds half of the primaries that serve hash slots, it escalates the state to **FAIL** and broadcasts a FAIL message to the whole cluster. Reports are only counted while they are recent — stale ones age out — so the agreement must form within a bounded window rather than accumulating over hours. Two consequences fall out of "majority of primaries": - **Three primaries minimum.** With two primaries, losing one leaves a single node, which is not a majority of two, so no FAIL can ever be agreed and no failover happens. Production clusters therefore run at least three primaries, ideally in three failure domains. - **Minority partitions are powerless.** The side of a partition holding fewer than half the primaries cannot declare FAIL, cannot authorise a promotion, and after `cluster-node-timeout` its primaries stop serving writes and report the cluster state as down. This is what keeps a partition from producing two authoritative primaries for the same slots. ## What FAIL triggers FAIL on a **primary that owns slots** is the precondition for failover: its replicas may now begin an election to take over the slots. FAIL on a replica, or on a node owning no slots, has no election consequence — it just marks the node unusable and stops it being counted. A node that was marked FAIL and comes back is not punished forever: if it is reachable again and it is a replica, or a primary that no longer serves any slot, the flag is cleared after a grace period so the node can rejoin normally. ## Timing: what a client actually experiences Failover is not instantaneous. Roughly: up to `cluster-node-timeout` for the first PFAIL, plus the gossip time for a majority to agree and broadcast FAIL, plus the replica's own election delay (a fixed few hundred milliseconds, a random component, and a rank-based component for lagging replicas), plus the time for clients to learn the new topology. With the 15 s default, a real-world unavailability window for the affected slots of the order of tens of seconds is normal. Requests to *other* slots are unaffected. ## Related settings worth naming - `cluster-node-timeout` — the detection threshold; also used as the base for replica eligibility (a replica whose link to its primary has been down far too long refuses to be promoted) and for how long an isolated primary keeps accepting writes. - `cluster-require-full-coverage` — when `yes` (the default), if any slot is uncovered the whole cluster refuses queries; when `no`, the covered slots keep serving. This decides whether an un-failed-over shard takes the entire cluster down with it. - `CLUSTER INFO` (`cluster_state`, `cluster_slots_ok`) and `CLUSTER NODES` (flags such as `pfail`, `fail`, `master`, `slave`) are how you observe all of this during an incident. ## Common misreading The most frequent error is to treat PFAIL and FAIL as the same thing, or to imagine that a client's timeout has anything to do with detection. Clients do not vote. Detection is entirely an internal, gossip-driven agreement among primaries, and only agreement — not any single node's view — starts a failover.

  • Why can a Redis Cluster with only two primaries never fail over automatically?
    Escalating PFAIL to FAIL requires PFAIL reports from a majority of the primaries that own slots. With two primaries, losing one leaves a single survivor, which is not a majority of two, so the FAIL state is never agreed and no replica is authorised to be promoted. Three primaries is the practical minimum, and they should sit in three separate failure domains so losing one domain still leaves a majority.
  • What happens to a primary that finds itself on the minority side of a network partition?
    It keeps serving briefly, then — once it has been unable to reach a majority of primaries for longer than cluster-node-timeout — it stops accepting writes and reports the cluster state as down, returning errors instead. That bounds how long it can accept writes that will be discarded when the majority side promotes a replica, which is the main mechanism limiting write loss during a partition.

saying these in an interview costs you the question

  • Treating PFAIL as sufficient to trigger a failover
  • Thinking client-side timeouts or client votes play any role in failure detection
  • Believing replicas participate in the majority count that escalates PFAIL to FAIL
  • Assuming failover completes within milliseconds of a node dying
  • Configuring different cluster-node-timeout values on different nodes

context

open as a page

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%

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.

open as a page

In Redis Cluster, a primary that owns a range of hash slots can keep replying OK to client writes for seconds after the cluster has already promoted one of its replicas for those same slots. Explain why that window exists, what happens to those writes when the old primary rejoins the cluster, and how the `cluster-node-timeout` setting sizes that window against the risk of unnecessary failovers.

level: seniorimportance: must knowfreq 50%

basics

~20 s

An isolated primary stops accepting writes only after failing to reach a majority of primaries for cluster-node-timeout, so it keeps serving its side meanwhile. On rejoin the promoted replica's higher configEpoch wins; the old primary demotes and discards those writes.

open as a page

What do the Redis settings `min-replicas-to-write` and `min-replicas-max-lag` actually protect against, and what do they fail to guarantee?

level: middleimportance: should knowfreq 35%

basics

~20 s

They make a primary refuse writes unless at least N replicas were acknowledging it within the last few seconds. That shrinks the window of writes that exist on only one node, but it is a check made before the write, not an acknowledgement after it — so accepted writes can still be lost.

open as a page

You are laying out a Redis Cluster across three availability zones. Which placement and replication decisions determine whether the cluster keeps serving — and keeps its writes — when an entire zone disappears?

level: principalimportance: should knowfreq 28%

basics

~20 s

Spread primaries so no zone holds half or more of them, since failover needs a majority of primaries; never place a replica in its primary's zone; keep spare replicas so a promoted node is not left bare; and decide explicitly whether uncovered slots should take the whole cluster down.

open as a page