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.
answer
- Isolated primary writes until majority unreachable for cluster-node-timeout
- Two owners for the same slots during the overlap
- Higher configEpoch wins on rejoin
- Old primary demotes to replica, resyncs, divergent writes gone
- Lower timeout = shorter window, more false failovers (fork/BGSAVE pauses)
basics
~20 sAn 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.
solid answer
~60 sA stranded primary has no instant signal that it lost. It keeps executing commands from clients that can still reach it, and only enters an error state — replying with errors instead of OK — once it has been unable to reach a **majority of primaries** for `cluster-node-timeout`. The majority side, meanwhile, has marked it failed and promoted a replica, so for part of that period two nodes believe they own the same hash slots and one of them is writing into a history nobody will keep. On rejoin, the old primary sees gossip claiming its slots with a **higher `configEpoch`**. Greater epoch wins: it reconfigures itself as a replica of the new owner and resyncs, **discarding** every write it accepted while isolated. There is no merge and no rollback notification. `cluster-node-timeout` is the knob: lower bounds the divergence window, but too low turns ordinary stalls — fork for BGSAVE, swap, cross-AZ jitter — into spurious failovers, each with its own resync and loss window. The small per-write acknowledgement gap and the `WAIT` / `min-replicas-to-write` mitigations are a separate, well-known axis.
code
text · 16 lines# --- during the partition, on the minority side ---
> SET order:77 paid
OK # accepted; the majority side has already promoted the replica
> SET order:78 paid
OK
# ...cluster-node-timeout elapses without reaching a majority of primaries
> SET order:79 paid
(error) CLUSTERDOWN The cluster is down
# --- after the partition heals, on the same node ---
> CLUSTER NODES
<old-primary-id> 10.0.1.5:6379 myself,slave <new-primary-id> 0 ... connected
<new-primary-id> 10.0.2.9:6379 master - 0 ... connected 0-5460
> GET order:77
(nil) # resynced from the new primary; the divergent writes are gonego deeper
Know the shape: a Redis primary cut off from the rest of the cluster can still say OK to writes for a while, and those writes disappear when it rejoins because the promoted replica is now the owner.
Be able to name the mechanism: the isolated primary stops serving only after cluster-node-timeout without reaching a majority of primaries, and on rejoin the higher configEpoch wins so it demotes to replica and resyncs.
Reason about the timeline concretely — overlap between the two owners, why a stalled-then-resumed primary shows the same symptom, why partial resync cannot save the surplus writes — and justify a cluster-node-timeout value from measured pause and RTT data rather than a default.
Frame it as a deliberate availability/consistency position: the window is bounded but real and unrecoverable, so decide what data is allowed to live only in Redis, make Redis-held state reconstructible, require idempotent retries, and use planned CLUSTER FAILOVER for anything you control the timing of.
## The setup Redis Cluster shards the keyspace into 16384 hash slots, each owned by exactly one primary, with zero or more replicas per shard. Nodes exchange gossip on a dedicated cluster bus port and agree on failures there; how a node moves from PFAIL to FAIL, and how a replica campaigns and wins an election, are covered separately in this topic. What matters here is the *other* side of that story: what the losing primary is doing while all of that happens. ## Why an isolated primary keeps accepting writes A Redis primary does not ask permission before serving a write. There is no lease it must renew, no quorum round on the write path — that is exactly why a single node sustains very high write throughput. Its only self-protection is a coarse, time-based check: if a primary has been unable to reach the **majority of primaries** for longer than `cluster-node-timeout`, it puts the cluster state into `fail` locally and starts replying with errors (`CLUSTERDOWN`) instead of serving reads and writes for its slots. So the sequence during a partition looks like this: 1. The partition happens. Clients on the minority side keep sending writes to primary **P**; P applies them and replies `+OK`. It has no idea anything is wrong. 2. On the majority side, P is gossiped as PFAIL then FAIL, and after the election delay one of its replicas **R** is promoted and takes ownership of the slots with a **new, higher `configEpoch`**. 3. Independently, P's own majority-unreachable timer expires and P stops accepting writes. Steps 2 and 3 are both roughly one `cluster-node-timeout` after the partition, but they are timed by different clocks and different starting events, and the election adds its own delay while P's cutoff does not. Real deployments also produce cleaner versions of the overlap: a primary that was **stalled** (a long fork/copy-on-write pause, swap, a hypervisor freeze) is failed over while unresponsive, then wakes up, briefly serves the clients still pointed at it, and only then discovers it is out of date. Clients with a cached topology map keep it company — a client only refreshes on a `MOVED` redirect or an error, and a stranded primary is issuing neither. The important framing for an interview: the exposure is **bounded by `cluster-node-timeout`, not by network round-trip time**. With the 15000 ms default, that is potentially fifteen seconds of writes that were acknowledged and are already doomed. ## What happens on rejoin: configEpoch decides Every primary carries a `configEpoch`, a monotonically increasing number that acts as the version of its slot-ownership claim. When R won its election it took a `configEpoch` strictly greater than P's. Slot-ownership conflicts in Redis Cluster are resolved by **the greater configEpoch wins** — this is the whole conflict-resolution rule, and it is why the cluster never ends up with two permanent histories. When the partition heals, P receives gossip and heartbeat packets advertising that R serves P's slots at a higher epoch. P does not argue and does not merge. It updates its slot map, **reconfigures itself as a replica of R**, and synchronizes. Because P wrote past the offset R was promoted at, its extra data cannot be reconciled by a partial resynchronization: it loads R's dataset and its own divergent writes are gone. No client is told; the application simply finds, later, that a write it saw succeed never existed. A manual `CLUSTER FAILOVER` is the deliberate contrast: the primary pauses clients, waits for the replica to match its offset, and hands over — a coordinated switch with no divergence. That is the tool for planned maintenance; the window described above is what an *unplanned* failover costs. ## Sizing cluster-node-timeout `cluster-node-timeout` (default 15000 ms) simultaneously controls: how long a node must be unreachable before peers suspect it, how long an isolated primary keeps writing, and — via `cluster-replica-validity-factor` — how stale a replica may be and still be electable. Lowering it directly shrinks the divergence window and shortens outages. The cost is **false failovers**: any stall longer than the timeout looks identical to death. Real sources of such stalls are ordinary — fork and copy-on-write pauses during `BGSAVE`/AOF rewrite on a large dataset, page-cache and swap pressure, noisy-neighbour or cross-zone network jitter, a slow `KEYS`-shaped command blocking the single-threaded event loop. Each false failover is not free: it triggers a promotion, a resync, a client topology refresh, and its own window of discarded writes. Setting it to a few hundred milliseconds converts a rare correctness problem into a frequent availability problem. The engineering answer is to measure: take the p99.9 of observed process pauses and inter-node round-trip time, and set the timeout comfortably above that — commonly 5000–10000 ms for a single-region cluster, higher when nodes span zones with jittery links. Then treat the residual window as real and design around it: keep money-like data in a store that acknowledges after durable quorum commit, make Redis-held state reconstructible, and make client retries idempotent (write an absolute value rather than `INCR`) so a retry after an ambiguous failure is safe.
- Why can't the rejoining primary keep its extra writes, or have them merged into the new primary?Redis Cluster has exactly one conflict-resolution rule for slot ownership: the claim carrying the greater configEpoch wins, and the loser adopts the winner's view wholesale. There is no vector clock, no per-key last-writer-wins merge, and no application hook to reconcile the two histories — Redis does not know whether two divergent values for a key are mergeable. Because the old primary wrote past the replication offset at which the replica was promoted, its surplus data cannot even be salvaged by a partial resynchronization; it resyncs and adopts the new primary's dataset.
- Would setting cluster-node-timeout to 500 ms be a reasonable way to make this window negligible?No — it trades a rare correctness problem for a frequent availability problem. Redis is single-threaded and routinely stalls longer than 500 ms for reasons that have nothing to do with failure: fork and copy-on-write pauses during BGSAVE or AOF rewrite on a large dataset, swap or page-cache pressure, an expensive O(N) command, cross-zone network jitter. Each spurious failover costs a promotion, a resync, a client topology refresh, and its own window of discarded writes. Pick the value from measured p99.9 pause and RTT data — usually somewhere in the 5000–10000 ms range for a single region.
- How does a planned failover avoid this window entirely?`CLUSTER FAILOVER`, issued on the replica you want promoted, is a coordinated handover rather than a detection-driven one. The current primary pauses client traffic, the replica waits until its replication offset matches, and only then does the role switch happen with a new configEpoch. No node is ever serving writes into a history that will be discarded, so it is the right tool for maintenance, host draining, or zone evacuation — the unplanned window is what you accept when the primary dies or is partitioned without warning.
A remote outpost commander who has lost radio contact keeps issuing orders on his own authority until the silence itself convinces him he is cut off. Headquarters has already appointed a successor with a later commission date. When contact returns, the later commission outranks him; every order he signed in the meantime is void, and he reports to the man who replaced him.
saying these in an interview costs you the question
- Claiming the old primary immediately knows it lost the slots and stops writing — it only finds out after failing to reach a majority for cluster-node-timeout, or on rejoin.
- Saying the two datasets are merged, or that the node with more writes / the higher replication offset wins — slot ownership is decided purely by the greater configEpoch.
- Believing the client gets an error or rollback for the writes that were discarded — they were acknowledged with +OK and are simply gone, silently.
- Assuming the exposure is a network round-trip or a few milliseconds; that is the separate acknowledgement gap. This window is measured in seconds and scales with cluster-node-timeout.
- Recommending a very small cluster-node-timeout as a free fix, ignoring fork/BGSAVE pauses and jitter that then cause repeated spurious failovers.