skip to content

In a search cluster that fans each query out to 100 shards, how do hedged requests tame tail latency caused by rare per-shard slowness?

level: seniorimportance: must knowfreq 55%

answer

  1. any-of-N probability
  2. one minus (1 − p) to the N
  3. duplicate only the late ones
  4. delay near the per-shard p95
  5. budget, cancel, and overload off-switch

basics

~20 s

If each shard is slow 1% of the time, a query that waits on 100 shards is slow about 63% of the time. A hedged request sends a backup copy to another replica once the first is late and uses whichever reply arrives first.

solid answer

~50 s

A query that waits for all `N` shards is slow if any one is slow: `1 − (1 − p)^N`. With `p = 1%` and `N = 100` that is about 63%, so most queries pay the per-shard p99. **Hedging** attacks that tail. Send each sub-request to one replica. If no reply arrives within a hedge delay of about the shard's p95, send a duplicate to a different replica, take whichever response comes first and cancel the other. With a p95 delay, at most about 5% of sub-requests are duplicated, yet the tail shrinks sharply when the slowness is replica-specific, such as a background pause, a noisy neighbour or a disk hiccup. Hedging does not help when the query itself is expensive on every replica. During cluster-wide overload it makes things worse, so a hedge budget and an overload signal should limit it.

code

pseudocode · 10 lines
pseudocode
function callShard(shard, req, hedgeDelay, deadline):
  first = send(shard.replica(0), req)
  if waitFor(first, hedgeDelay):          // on time: no hedge
    return first.result
  pending = [first]
  if hedgeBudget.tryAcquire():            // cap hedges during overload
    pending.add(send(shard.replica(1), req))
  winner = waitAny(pending, until = deadline)
  cancelAllExcept(pending, winner)        // winner null: cancel all
  return winner == null ? TIMED_OUT : winner.result

go deeper

for a junior

Recall that a query waiting on many shards is as slow as its slowest shard, and that a hedged request sends a late sub-request again to another replica and keeps the first reply.

for a middle

Explain the any-of-N formula and be able to compute it for a few shard counts. Walk through the hedge steps: delay, backup to a different replica, first reply wins, loser cancelled.

for a senior

Show operating judgment: choose the hedge delay from the latency distribution, cap hedges with a budget, switch them off under overload, and recognise query-intrinsic slowness that hedging cannot fix.

for a principal

Weigh hedging against tied requests, replica selection, deadlines with partial results and shard-count changes. Decide which layers a latency SLO needs and what extra capacity each one costs.

## Tail-latency amplification A single shard is fast most of the time, but every server has occasional slow requests. Causes include background compaction, a memory-management pause, a noisy neighbour, a cold cache or a retransmitted packet. Suppose each shard independently takes longer than 1 s on `p = 1%` of requests. A scatter-gather query that waits for **all** `N` shards is slow if **any** of them is slow: `P(query slow) = 1 − (1 − p)^N` | Shards waited on (N) | Share of queries slower than 1 s | |---|---| | 1 | 1.0% | | 10 | 9.6% | | 50 | 39.5% | | 100 | 63.4% | | 1,000 | effectively 100% | At 100 shards, most queries (about 63%) pay the per-shard p99. Improving the *average* shard latency does not help. The fix has to target the tail directly, or stop waiting for it. These figures assume shards slow down independently. Correlated slowness, such as cluster-wide overload, behaves differently and is covered below. ## Hedged requests A **hedged request** duplicates only the sub-requests that are already late: 1. Send the sub-request to one replica of the shard. 2. Wait for a **hedge delay**, typically about that shard's p95 latency. 3. If no reply has arrived, send the same sub-request to a **different replica**. 4. Use whichever response arrives first and **cancel** the other copy. 5. Enforce a **hedge budget**, for example no more than a few percent of sub-requests hedged, so that a widespread slowdown cannot double traffic. ```pseudocode function callShard(shard, req, hedgeDelay, deadline): first = send(shard.replica(0), req) if waitFor(first, hedgeDelay): return first.result pending = [first] if hedgeBudget.tryAcquire(): pending.add(send(shard.replica(1), req)) winner = waitAny(pending, until = deadline) cancelAllExcept(pending, winner) return winner == null ? TIMED_OUT : winner.result ``` ## What hedging costs - **Extra load.** With the delay at the p95, only about 5% of sub-requests are still pending when the timer fires, so the extra sub-requests are at most about 5%. The real figure is lower when cancelled losers stop early. - **Latency bound.** When the first replica stalls, a hedged sub-request typically finishes within the hedge delay plus a normal response time from the second replica. - **Duplicate work.** Both replicas may run the query until the cancellation reaches one of them. Search reads have no side effects, so duplicates are safe, but they still consume CPU. Sending two copies of every sub-request immediately would also cut the tail, but it doubles the load for a benefit that comes almost entirely from the few slow cases. ## When hedging backfires or does not help - **Query-intrinsic slowness.** An expensive query, such as a huge wildcard expansion or a very common term, is slow on every replica. The backup is just as slow and only adds load. - **Cluster-wide overload.** When the slowness comes from saturation, hedges push more traffic into queues that are already full and deepen the incident. The budget and an overload signal (queue depth, rejection rate) should switch hedging off. - **Correlated replicas.** If two replicas share a disk, a host or a network path, the backup inherits the stall. - **No cancellation.** If the losing copy is not cancelled, the extra work runs to completion and hedging costs more than planned. ## Related techniques | Technique | Idea | Trade-off | |---|---|---| | Tied requests | send to two replicas at once; whichever starts executing first tells the other to drop its copy | almost no duplicate work, but replicas must message each other | | Latency-aware replica selection | route to replicas with low recent latency and short queues | avoids known-slow replicas, but can herd traffic onto a few | | Deadline with partial results | stop waiting at a deadline and answer from the shards that replied | bounded latency, possibly incomplete answers | | Fewer, larger shards | reduce `N` | each shard is slower and less parallel | | Result caching | answer repeated queries with no fan-out | helps only popular queries and must be invalidated when the index refreshes | In practice these techniques are layered: replica selection first, hedging within a budget, and a deadline as the final backstop.

  • In a fanned-out search, why wait for about the p95 before hedging instead of sending two copies immediately?
    Most of the benefit comes from the few requests that are actually slow. Waiting until the p95 means only about 5% of sub-requests get a duplicate, so the extra load stays small while the worst stalls are still cut short. Sending two copies of everything doubles the load for little extra gain. Tied requests are a middle ground: both copies go out at once, and each replica cancels its sibling when it starts executing.
  • When can hedging in a search cluster make an incident worse?
    Hedging hurts when the slowness comes from overload rather than from one unlucky replica. In that case every hedge adds work to queues that are already saturated, which increases latency further and triggers still more hedges. Guard against this with a hedge budget capped at a small share of traffic, prompt cancellation of losers, and an automatic off-switch driven by queue depth or rejection rates.
  • Besides hedging, what bounds the latency of a query that fans out to many search shards?
    Latency-aware replica selection steers sub-requests away from slow replicas before they are sent. A per-query deadline caps the wait, and the coordinator answers from the shards that replied. Reducing the shard count lowers the any-of-N exposure, and a result cache answers repeated queries without any fan-out. These techniques are layered, with the deadline as the final backstop.

saying these in an interview costs you the question

  • A shard's p99 affects only about 1% of fanned-out queries.
  • Hedging doubles cluster load because every request is sent twice.
  • Hedging fixes queries that are expensive on every replica.
  • Raising the timeout is the right fix for fan-out tail latency.
  • Hedged requests are safe to send even when the whole cluster is overloaded.