skip to content

How does a SolrCloud shard elect its leader, and why can a shard end up with no leader?

level: seniorimportance: should knowfreq 45%

answer

  1. ZooKeeper provides the ordering, not a vote
  2. The znodes are ephemeral and sequential
  3. Each candidate watches only its predecessor
  4. A counter tracks who is up to date
  5. Some replica types are never candidates

basics

~20 s

Leader-eligible replicas of a shard queue up on ephemeral sequential ZooKeeper nodes; the lowest sequence wins and publishes itself as leader. A shard is left leaderless when no eligible replica is available or up to date — for example only PULL replicas survive, or every replica is down.

solid answer

~60 s

Each shard runs its own election. Leader-eligible replicas (NRT and TLOG — never PULL) create **ephemeral sequential** znodes under the shard's election path in ZooKeeper. The lowest sequence number becomes leader and publishes itself; each other candidate watches only the node immediately ahead of it, so a leader's departure wakes one candidate rather than a thundering herd. Because the znodes are ephemeral, a crashed node's candidacy disappears with its ZooKeeper session. Solr does not elect by replica vote or majority — ZooKeeper provides the ordering. What Solr adds is **shard terms**: a per-shard counter recorded in ZooKeeper that tracks which replicas are current. A replica that failed to accept an update is left behind on terms, cannot be elected, and must recover first. A shard has no leader when nothing eligible and current is available: only PULL replicas survive, all replicas were down and the first one back is still waiting out `leaderVoteWait` for a possibly-better peer, or ZooKeeper quorum is lost so no election can run at all.

code

bash · 2 lines
bash
# Inspect which replica currently leads each shard
curl "http://localhost:8983/solr/admin/collections?action=CLUSTERSTATUS&collection=products"

go deeper

for a junior

Know that leadership is per shard, that ZooKeeper decides it, and that the leader coordinates writes while any active replica can serve queries.

for a middle

Explain the ephemeral sequential znode queue, the predecessor-watch pattern that avoids a herd, and which replica types are eligible candidates.

for a senior

Diagnose a leaderless shard in production: distinguish only-PULL-survivors, an open leaderVoteWait window, stale shard terms and lost ZooKeeper quorum, and know why FORCELEADER is a last resort.

for a principal

Set the standards that prevent it — minimum leader-eligible replicas per shard, rolling-restart policy, leadership rebalancing after maintenance, and an explicit position on write loss versus write availability.

## The election mechanism Leadership in SolrCloud is per **shard**, not per collection and not per node. A single node can be leader of one shard and a follower for another. When a leader-eligible replica starts up (or the current leader disappears), it creates an **ephemeral sequential** znode under that shard's election path in ZooKeeper. ZooKeeper hands out monotonically increasing sequence numbers, producing a queue. The replica holding the lowest sequence number becomes leader and publishes its identity under the collection's leader path so every node can see it. Two properties of ZooKeeper do the heavy lifting: - **Ephemeral** — the znode vanishes when that replica's session ends, so a crash or a long GC pause that expires the session automatically withdraws the candidacy. - **Sequential** — the ordering is decided by ZooKeeper, so there is no vote to run among replicas and no split-brain from two replicas both concluding they won. Each waiting candidate watches the node immediately ahead of it in the queue rather than the leader itself. When a node leaves, exactly one watcher fires. That is the standard herd-avoidance pattern and it matters on clusters with many shards. ## What the leader actually does The leader is the shard's write coordinator. Every update for that shard is routed to it; it assigns the document's `_version_`, writes to its own transaction log and index, and forwards the update to the shard's other replicas in parallel, waiting for their acknowledgements. Queries, by contrast, may be served by any active replica — leadership is about the write path, not read capacity. If a replica fails to apply a forwarded update, the leader records that divergence through **shard terms**, a per-shard counter in ZooKeeper. Replicas that keep up share the current term; one that missed an update is left behind. Two consequences follow: a lagging replica knows it must recover, and it is excluded from becoming leader while it is behind. This term mechanism replaced Solr's older leader-initiated-recovery approach and is what prevents a stale replica from winning an election and quietly rolling back acknowledged writes. ## Recovery after election When a new leader takes over, other replicas compare their term and version state against it. A replica that is only slightly behind replays from its transaction log — a *PeerSync*-style catch-up. One that is too far behind gives up on incremental catch-up and pulls a full index copy from the leader, which is expensive in disk and network but always correct. A TLOG replica elected leader replays what it needs and then switches into local indexing mode for as long as it holds leadership. ## Why a shard can be leaderless 1. **No eligible replica alive.** PULL replicas can never be elected — they keep no transaction log, so they cannot demonstrate they hold every acknowledged update. A shard with one NRT replica and three PULL replicas has exactly one leader candidate; lose it and the shard cannot take writes. 2. **All replicas were down.** When the whole shard restarts, the first replica back does not seize leadership immediately: it waits (`leaderVoteWait`) for peers to appear, so that a stale replica does not become the authority while a more current one is thirty seconds from starting. During that window the shard has no leader, by design. 3. **Everyone is behind.** If shard terms say no live replica is current, none can be safely elected. 4. **No ZooKeeper quorum.** Election is implemented with ZooKeeper writes; without a majority ensemble, no election can run at all. Symptoms: updates to that shard fail while other shards keep working, CLUSTERSTATUS shows the shard with no `leader` flag on any replica, and logs mention waiting for a leader. ## Operating it - Keep **at least two leader-eligible replicas per shard**. This is the single most common design mistake in PULL-heavy layouts. - Prefer a rolling restart to a full shard restart, so leadership can move rather than lapse. - The Collections API `FORCELEADER` action exists for a shard that is genuinely stuck with no leader and no path to one. It is a last resort: it can promote a replica that is not fully up to date, which risks losing acknowledged updates. Understand why the shard is stuck before reaching for it. - `REBALANCELEADERS` can redistribute leadership after a restart concentrates leaders on one node — worth knowing when one node's CPU is pinned while its peers idle. - Watch for replicas cycling through `recovering` — repeated full index copies usually mean a replica that keeps falling behind rather than a one-off failure.

  • What are shard terms and what problem do they solve?
    A shard term is a per-shard counter in ZooKeeper tracking which replicas are current. When a replica fails to apply an update the leader forwarded, the leader advances the term for the healthy replicas, leaving the failed one behind. That replica then knows it must recover and, crucially, is excluded from leader election while it is behind — preventing a stale replica from being elected and silently discarding acknowledged writes.
  • Why does each election candidate watch only the znode ahead of it instead of watching the leader?
    To avoid a thundering herd. If every candidate watched the leader's znode, a single leader loss would wake all of them at once, and ZooKeeper would field a burst of simultaneous reads and writes — multiplied by every shard in the cluster during a node failure. Watching only the immediate predecessor means one departure wakes exactly one candidate, keeping election traffic proportional to the change.
  • When is it acceptable to use the FORCELEADER action?
    Only when a shard is genuinely stuck with no leader and no route back to one — for instance the last fully current replica's data is gone for good. FORCELEADER can promote a replica that is not fully up to date, so acknowledged updates may be lost. Before using it, verify why no election is completing: quorum loss, only PULL replicas alive, and a leaderVoteWait window still open all look similar but need different responses.
  • Why doesn't the first replica back after a full shard outage become leader immediately?
    Because it may be the stale one. Solr makes it wait (governed by leaderVoteWait) for its peers to appear, so a replica that missed the last minutes of indexing does not become the authority while a more current one is still starting. The shard is briefly leaderless on purpose — that pause trades a short write outage for not losing acknowledged updates.

saying these in an interview costs you the question

  • Says replicas vote for a leader by majority
  • Claims PULL replicas can be promoted in an emergency
  • Thinks the leader also handles all queries for the shard
  • Reaches for FORCELEADER before diagnosing the cause
  • Believes leadership is per collection rather than per shard

context