skip to content

Explain lease-based leadership in an automated relational failover setup such as Patroni backed by etcd: what the leader key and its TTL are, what each agent does on every loop, and how this arrangement stops two primaries from existing.

level: seniorimportance: should knowfreq 34%

answer

  1. leader key + TTL in etcd
  2. renew or self-demote
  3. CAS create only after expiry
  4. loop_wait well below ttl
  5. DCS quorum loss = read-only by design

basics

~20 s

Leadership is a short-lived key in a quorum-backed store (etcd) that the primary's agent must keep renewing. If it cannot renew before the TTL expires, it demotes its own database; only after the key expires can another agent create it and promote. One key, one writer.

solid answer

~60 s

The database engines themselves have no notion of cluster leadership, so an agent per node (Patroni) plus a consensus store (etcd, Consul, ZooKeeper) supplies one. Leadership is represented as a single key in the store holding the current leader's name, with a short TTL, typically tens of seconds. Each agent loops on a fixed interval: check its local database, read cluster state from the store, and act. The leader's agent renews the key. Every follower's agent reads the key; if it exists it stays a replica pointed at that leader. The safety comes from two halves. If the leader agent cannot renew, because the store is unreachable, it lost quorum, or its own database is unhealthy, it demotes the local database to read-only itself, before the TTL elapses. Meanwhile no other agent can take the key until it actually expires, and creating it requires a quorum write through the store's consensus protocol, so only one candidate wins. The demote-before-expiry ordering is what makes a second promotion safe. The key parameters are the loop interval, the TTL, and the store's own timeouts; the TTL must exceed the loop interval by a comfortable margin or you get spurious failovers.

code

text · 5 lines
text
/service/mycluster/leader      = "pg-node-1"   (ttl 30s, refreshed every 10s)
/service/mycluster/members/pg-node-1 = {role: master, state: running, xlog_location: 0/7A3C1F0}
/service/mycluster/members/pg-node-2 = {role: replica, state: running, xlog_location: 0/7A3C1F0}
/service/mycluster/members/pg-node-3 = {role: replica, state: running, xlog_location: 0/7A3B980}
/service/mycluster/config      = {ttl: 30, loop_wait: 10, maximum_lag_on_failover: 1048576}

go deeper

for a junior

Know the shape: an agent next to each database plus a shared store holding a short-lived leader key that must be renewed.

for a middle

Walk the agent loop and state the two-sided safety property: renew or demote, and no one takes the key before it expires.

for a senior

Reason about the timing budget (loop interval vs TTL vs store latency), watchdog pairing, and the deliberate read-only behaviour on store quorum loss.

for a principal

Treat the consensus store as a first-class availability dependency: size and place it, set the TTL from measured latency distributions, and decide what the system does when it is gone.

## The pieces PostgreSQL and MySQL know how to be a primary or a replica but have no built-in concept of which node in a cluster should be the primary. Automated failover therefore adds two components: - A **distributed configuration store (DCS)** with consensus and quorum semantics: etcd, Consul, or ZooKeeper. It provides linearizable, compare-and-swap style writes and keys with a time to live. - An **agent per database node** (Patroni is the canonical PostgreSQL example) that owns the local database process and reads/writes cluster state in the DCS. All cluster state lives in the DCS: cluster configuration, each member's status and replay position, and one leader key whose value is the current leader's name. ## The lease The leader key is a lease, not a flag. It carries a TTL (Patroni's ttl, commonly 30s). Holding leadership means having successfully written that key recently enough that it has not expired. On every loop (Patroni's loop_wait, commonly 10s) each agent does roughly: 1. Query the local database: is it up, is it primary or replica, what is its replay position? 2. Read cluster state from the DCS and publish its own member status. 3. If it is the leader: verify the local database is a healthy primary, then refresh the leader key with a fresh TTL. 4. If it is not the leader and the leader key exists: ensure the local database is running as a replica following the named leader. 5. If the leader key is absent: evaluate whether it is a legitimate candidate (healthy, not lagging beyond the configured threshold, no better candidate) and, if so, attempt to create the key with a compare-and-swap that succeeds only if the key still does not exist. Because step 5 is a quorum write mediated by consensus, exactly one agent can win the race, even if several attempt it in the same millisecond. ## Why two primaries cannot both persist The dangerous window is: old leader alive but partitioned, new leader promoted. Lease semantics close it from both ends. - **From the old leader's side**: it can only remain primary while it keeps renewing. If it cannot reach a quorum of the DCS, whatever the reason (network partition, DCS quorum loss, its own hang), the renewal fails and the agent demotes the local database to read-only or shuts it down. Crucially it does this proactively, on failure to renew, not on hearing about a new leader, because in a partition it hears nothing. - **From the candidate's side**: no agent may create the leader key until the existing lease has expired in the DCS's own clock, and the DCS itself only accepts the write with a quorum. So promotion can only start after the point by which the old leader must already have given up. The correctness argument is therefore a timing one: demote-on-failed-renewal must complete before the key expires and someone else promotes. That is why the TTL must be comfortably larger than the loop interval and the DCS request timeout, and why a hardware watchdog is recommended: if the agent process is wedged and cannot even demote, the watchdog resets the host within the same budget. ## Tuning the numbers Roughly, the TTL bounds the worst-case failover time (nobody promotes before the lease expires) and also bounds how long a partitioned old primary might still be writing. Shrinking it makes failover faster but makes transient network blips or a slow DCS look like failures, producing spurious failovers and needless promotions. The relationship loop_wait plus a retry budget must fit inside the TTL; managers refuse configurations that violate it. Because the DCS is now on the critical path, its own latency and health matter: a saturated or disk-slow etcd cluster can trigger database failovers even though every database is fine. ## Consequences worth naming - Losing DCS quorum means the cluster goes read-only by design, even if all databases are healthy. That surprises people; it is the correct trade. - The DCS must be sized and placed like any quorum system: odd membership, independent failure domains, and not co-located so that one rack takes both database and store. - Leadership in the store is only half the story: client traffic must follow the key, which is why these setups pair with a routing layer that reads the same state. ## Interview framing Name the two components, describe the key as a TTL lease, walk one agent loop, and then give the safety argument in one sentence: the holder demotes itself on failed renewal before anyone else is allowed to take an expired lease, and taking it requires a quorum write.

  • What happens to a perfectly healthy three-node database cluster if the etcd cluster loses quorum?
    The leader's agent can no longer refresh the leader key, so it demotes its database to read-only, and no candidate can create the key either, so the cluster serves reads only until the store recovers. This is intentional: without a quorum store there is no safe way to prove single-writer status. It also means the store must be operated with the same care as the databases, since its availability is now the cluster's write availability.
  • Why not simply set the TTL to two seconds for faster failover?
    Because the TTL must exceed the agent loop interval plus retry budget plus normal store latency; below that, ordinary jitter, a garbage collection pause, or a slow etcd disk looks like a dead leader and triggers a promotion. Each spurious failover costs a connection storm, a rewind or rebuild of the demoted node, and aborted in-flight transactions. The usual practice is to pick the smallest TTL that survives your observed p99 latencies with margin.

saying these in an interview costs you the question

  • Thinking the new leader is promoted the moment the old one stops answering, rather than after lease expiry
  • Forgetting that the old leader demotes itself on failed renewal, and assuming only the candidate side matters
  • Setting the TTL close to or below the agent loop interval
  • Treating the consensus store as optional infrastructure that does not need quorum sizing or monitoring
  • Assuming the database engine itself arbitrates leadership

context