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?
answer
- Cluster bus = data port + 10000, gossip in PING/PONG
- PFAIL = local opinion after cluster-node-timeout
- FAIL = majority of slot-owning primaries agree, then broadcast
- Needs ≥3 primaries for any majority to exist
- FAIL on a primary unlocks replica election
basics
~20 sNodes 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 sEvery 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> 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
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.
Explain the gossip mechanism, the majority-of-primaries requirement, and why three primaries is the minimum for automatic failover.
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.
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