skip to content

How would you configure a MongoDB replica set to meet a strict limit on write unavailability during failover?

level: principalimportance: should knowfreq 30%

answer

  1. Break the outage into phases first
  2. Detection is usually the biggest slice
  3. One lever lives in the driver, not the server
  4. Voter placement bounds the voting round trip
  5. Rehearse with a controlled step-down

basics

~20 s

Break the outage into detection, voting, catch-up and client rediscovery, then attack each. Size the voters odd and keep them on a low-latency network, avoid arbiters, use retryable writes so a failover is latency rather than an error, and rehearse with rs.stepDown() to measure the real number.

solid answer

~50 s

Start by decomposing the window rather than tuning one setting. Detection is bounded by `electionTimeoutMillis`, default `10000`, and normally dominates. Voting costs roughly one round trip to the slowest member of the voting majority. Catch-up is the new primary applying entries it was missing, bounded by `catchUpTimeoutMillis`. Client rediscovery is the driver noticing the new primary and re-selecting a server. Then act on each: keep an odd number of voters and place them where round trips between them are short, since the majority's slowest member sets the voting cost; avoid arbiters, whose degraded mode is far worse than a third data-bearing member's; use `priority` to steer the primary near the clients but remember a priority change is itself a failover; enable retryable writes so one failover becomes added latency instead of an application error; and use `w: "majority"` so the failover cannot silently discard acknowledged writes. Finally, rehearse with `rs.stepDown()` under production-shaped load and measure the client-visible outage — that number, not the config, is your budget.

code

javascript · 3 lines
javascript
// Controlled rehearsal: hand over, then measure from the client
rs.stepDown(60, 10)   // stepDownSecs, secondaryCatchUpPeriodSecs
rs.status().members.map(m => ({ name: m.name, state: m.stateStr }))

go deeper

for a junior

Recall that a replica set is briefly unable to accept writes while it elects a new primary, and that the wait before an election starts is governed by electionTimeoutMillis.

for a middle

Be able to name the phases — detection, voting, catch-up, client rediscovery — and say which setting bounds each rather than attributing the whole window to one knob.

for a senior

Demonstrate that you have measured a real failover from the client side, know that retryable writes change the user-visible outcome, and can explain why an aggressive timeout trades one outage for several.

for a principal

Own the availability contract: what window the business actually needs, what each compression costs in stability and money, where voters sit, and what the application does during the window when the number cannot go lower.

## Reframe the question first "Make failover faster" is usually answered with "lower `electionTimeoutMillis`", which is a partial answer that trades one problem for another. The productive framing is: **what is the client-visible write-unavailable window made of, and which parts can I actually shrink?** There are four sequential phases. **1. Detection.** Members heartbeat each other on `heartbeatIntervalMillis` (default `2000`). A secondary that goes `electionTimeoutMillis` (default `10000`) without successful contact with the primary concludes the primary is gone. On a default configuration this phase is most of the outage. **2. Voting.** The candidate solicits votes and needs a strict majority of voting members. This is roughly one round trip, and it is bounded not by the average latency among voters but by the latency to the slowest member of whichever majority answers. Where the voters sit therefore directly sets this cost. **3. Catch-up.** The elected primary attempts to apply oplog entries other members hold and it does not, before accepting writes, so that as little history as possible is discarded. `catchUpTimeoutMillis` bounds this phase. A well-replicated set spends almost no time here; a set with a lagging secondary can spend a lot. **4. Client rediscovery.** Drivers monitor topology on their own interval. Until the driver learns of the new primary, operations wait in server selection and eventually fail against the driver's server-selection timeout. This phase is invisible from the server side and is frequently the part teams forget. ## Levers, in the order they usually pay **Enable retryable writes.** This is the highest-leverage change and it is client-side. With retryable writes, a write that fails because the primary went away is retried once after the driver re-selects a server. A single failover then shows up as one slow request rather than an error surfaced to the user. It does not shorten the server-side window, but it changes what the window *means* to the application — which is what the availability target is actually about. **Place the voters for fast agreement.** Voting latency is set by the slowest member of the majority. Voters spread across distant regions make every election slow and also make aggressive detection timeouts unsafe. Keep the voting members on a low-latency network with each other; put distant copies in as `priority: 0` (and, if you want them out of the arithmetic entirely, `votes: 0`) members. **Use an odd number of voters.** Even counts add a voter without adding fault tolerance and create a symmetric-split shape in which nobody has a majority and the set is write-unavailable indefinitely — an unbounded outage, which no timeout tuning can fix. **Avoid arbiters.** An arbiter buys a vote and nothing else. In a primary-secondary-arbiter set, losing the single secondary leaves majority write concern unsatisfiable and the majority commit point frozen, which degrades the surviving primary over time. A third data-bearing member costs a machine and removes that entire class of incident. **Keep secondaries caught up.** Replication lag lengthens the catch-up phase and widens the set of writes exposed to rollback. Lag is worth monitoring specifically as a failover-readiness signal, not only as a read-staleness signal. **Tune `electionTimeoutMillis` last, and only with evidence.** Lowering it shortens detection but makes GC pauses, disk stalls and latency spikes look like a dead primary. Each false positive costs another outage window and risks rolling back writes acknowledged only by the old primary. Lower it only where you have measured the tail latency between voters and have headroom. **Use `w: "majority"` on writes that matter.** This does not shorten the outage. It ensures the outage does not also silently lose acknowledged writes, which is usually the more expensive failure. ## Steering where the primary lands `members[n].priority` biases which member the set converges on, so you can keep the primary in the region holding most of the clients. Two cautions. A priority change is not passive — a higher-priority secondary calls a takeover election as soon as it is caught up, producing a real failover at the moment you run the reconfig. And a flapping high-priority node will pull the primary back to itself repeatedly, turning one bad machine into a series of outages. Priority steers; it does not guarantee. ## Measure, don't model The only credible number is a measured one. Rehearse with `rs.stepDown()` on the primary under production-shaped load, and instrument the *client*: time from step-down to the first successful write from the application, not from the server logs. A planned step-down is the optimistic case because it includes a catch-up courtesy period; also rehearse the pessimistic case by hard-stopping a primary in a non-production environment, where detection runs the full timeout. Run this regularly rather than once. Voter placement drifts, drivers get upgraded, lag profiles change with load. A failover budget that was verified a year ago is a claim, not a measurement. ## What to say about the trade-off The honest summary is that you can compress the outage to a few seconds with careful placement and a modestly reduced timeout, and you can make most of that invisible to users with retryable writes — but you cannot drive it to zero, and pushing detection too hard converts a rare outage into frequent ones. If the requirement is genuinely zero write unavailability, the answer is not replica set tuning; it is a change in what the application does during the window, such as queueing or degrading gracefully.

  • Why is enabling retryable writes often more valuable than lowering electionTimeoutMillis?
    Lowering the timeout shortens the server-side window but raises the rate of spurious failovers, so you may trade one long outage for several short ones. Retryable writes change the outcome instead of the duration: the driver re-selects a server and retries the failed write once, so a single failover surfaces as extra latency rather than an error. It also costs nothing in stability.
  • How would you actually measure your failover budget rather than estimate it?
    Run a controlled `rs.stepDown()` on the primary while the system carries production-shaped load, and instrument from the client side — time from step-down to the first successful application write, not server log timestamps. Repeat with a hard primary stop in a non-production environment to capture the pessimistic case where detection runs the full election timeout. Re-run it on a schedule, since placement, drivers and lag profiles drift.
  • A stakeholder asks for zero write unavailability during failover. How do you respond?
    Explain that a replica set has one writable primary by design, so some window always exists between losing it and electing another. You can compress it and make it invisible to most users with retryable writes, but not eliminate it. If the requirement is real, it has to be met in the application — queue writes for the duration, degrade to a read-only experience, or accept eventual application of the write.
  • Which phase of a failover does replication lag affect, and why does it matter beyond speed?
    Catch-up. The elected primary applies oplog entries other members hold before accepting writes, so a lagging set spends longer in that phase. Beyond duration, lag widens the set of writes sitting past the majority commit point, so a failover during high lag discards more acknowledged-but-unreplicated writes as rollback. Lag is a failover-readiness signal, not just a read-staleness one.

saying these in an interview costs you the question

  • Answers only with 'lower electionTimeoutMillis'
  • Ignores the client-side rediscovery phase entirely
  • Adds voters believing more members means faster failover
  • Quotes a failover time that was never measured from the client
  • Treats a priority change as a passive configuration edit

context