skip to content

After sharding, p99 read latency got worse and every shard is busy — how do you confirm broadcast queries are the cause?

level: seniorimportance: must knowfreq 58%

answer

  1. Diagnose per query shape, not per server
  2. The winning plan names the participants
  3. Most participants return nothing at all
  4. Latency is the worst of N, not the average
  5. More shards makes this particular problem worse

basics

~20 s

Run explain() through mongos on the hot query shapes: a SHARD_MERGE stage with every shard named in the winning plan means scatter-gather. Compare per-shard keysExamined against nReturned, then fix the filter to carry the shard key rather than adding shards.

solid answer

~50 s

Start from the query shapes, not the hardware. Run `explain("executionStats")` through `mongos` for each hot read. A **`SINGLE_SHARD`** top-level stage means the query was targeted; **`SHARD_MERGE`** with all shards listed in the `shards` array means it broadcast. In `executionStats` you get per-shard `nReturned`, `keysExamined` and `docsExamined` — a broadcast typically shows most shards returning zero while still paying a scan. The latency story follows: a fan-out query's response time is the **maximum** over N shards plus the merge, so you are sampling the tail of the shard population on every request, and one slow shard sets p99 for everything. Meanwhile per-shard load is proportional to total query rate rather than to that shard's share of data, so adding shards makes it worse, not better. The fix is to make the hot filters carry the shard key — reshape the query, or carry the key on the access path — not to buy capacity.

code

javascript · 5 lines
javascript
// Run through mongos, per hot query shape
db.orders.find({ status: "OPEN" }).explain("executionStats")
// queryPlanner.winningPlan.stage: "SHARD_MERGE"  -> broadcast
// queryPlanner.winningPlan.shards: [ ... all shards ... ]
// executionStats: per-shard nReturned / keysExamined / docsExamined

go deeper

for a junior

Know that explain() run through mongos tells you how many shards a query used, and that a query hitting every shard is the expensive kind.

for a middle

Be able to read a sharded explain end to end: the winning plan's stage, the shards array, and the per-shard nReturned versus keysExamined that exposes wasted work.

for a senior

Drive the diagnosis: quantify shards touched per query weighted by call rate, explain why the tail degrades before the median, and fix the filter rather than the capacity.

for a principal

Own the guardrail. Decide how fan-out is measured and reported release over release, and when a workload's shape justifies resharding or a separate serving path instead of more hardware.

## The symptom and the wrong instinct The classic report is: "we sharded to scale, and now the same queries are slower and every shard is at the same CPU." The instinct is to add a shard. That is the one action guaranteed not to help, because if the queries broadcast, each new shard adds a participant to every request rather than taking a slice of the traffic. ## Step 1: confirm with explain, per query shape Routing is a per-query property, so diagnose per shape. Through `mongos`: ```js db.orders.find({ status: "OPEN" }).explain("executionStats") ``` What to read: - **Top-level winning-plan stage.** `SINGLE_SHARD` means one shard was contacted. `SHARD_MERGE` means several were and the router merged their output. - **The `shards` array.** It names the participants. If it lists every shard in the cluster, the query is a full broadcast. - **Per-shard `executionStats`.** For each shard you see `nReturned`, `keysExamined`, `docsExamined` and its plan. The tell of a wasteful broadcast is many shards with `nReturned: 0` that nonetheless examined keys or, worse, ran a `COLLSCAN`. - **`SHARDING_FILTER`.** Shard plans include this stage; it discards orphaned documents that a shard still holds from an interrupted migration, so the router never returns them twice. Do this for the top handful of shapes by call volume. Aggregations get the same treatment — the sharded explain shows which shards ran the pipeline and how it was split. ## Step 2: quantify the fan-out Two numbers make the case to anyone: 1. **Shards touched per query.** Weight each hot shape by its call rate. If 80% of calls touch all 20 shards, the cluster is doing roughly 20× the logical work of a single node for those calls. 2. **Per-shard request rate.** On a targeted workload each shard sees about `total_qps / N`. On a broadcast workload each shard sees `total_qps`. Comparing the observed per-shard operation counters against total application throughput settles the question without any explain at all. ## Step 3: understand why p99 in particular degrades Fan-out queries have a structural tail problem. The merge cannot produce its first result until every participating shard has answered, so each request's latency is the **maximum** of N shard latencies, not the mean. If any single shard is momentarily slow — a checkpoint, a competing scan, an election, a noisy neighbour — every in-flight broadcast query inherits that latency. With N shards you draw N samples per request, so your p99 comes from a distribution far to the right of any individual shard's p99. This is why the median can look acceptable while p99 falls apart, and why the effect worsens as you add shards. Secondary amplification: each broadcast consumes a connection and a worker thread on every shard, and the router holds N cursors per query, so connection-pool pressure and cursor bookkeeping grow with the same factor. ## Step 4: fix the routing, not the capacity In rough order of preference: - **Put the shard key in the filter.** Very often the value is available at the call site but was never passed down — a tenant id, an account id, a user id sitting in the request context but not in the repository method's arguments. Threading it through turns a 20-shard query into a 1-shard query with no data change. - **Reshape the query, not the cluster.** If a screen fetches by a secondary attribute, ask whether the natural entry point is the key-bearing entity instead. - **Serve the shape from a purpose-built collection.** A materialised summary keyed so that the reads become targeted moves the fan-out off the request path onto a background job. - **Accept and contain it.** Some queries genuinely cannot be targeted. Keep them off the latency-critical path, bound their result size, and make sure an index supports them on every shard so each participant's contribution is a short index walk rather than a scan. Resharding the collection on a different key is possible and is the heavier hammer; that decision belongs with shard-key design, and it is worth reaching only when the whole hot workload shares a key that the current one lacks. ## Step 5: keep watching Make fan-out visible in normal operations: track shards-touched for hot shapes after each release, and watch per-shard operation counters relative to application throughput. Broadcast creep is a silent regression — one new filter shipped by one team can turn a well-targeted endpoint into a cluster-wide one without any error, any alert, or any change in the median.

  • Why does adding shards make a broadcast-heavy workload slower rather than faster?
    Each additional shard becomes another participant in every request, so per-query work and connection use grow while nothing takes traffic away from the others. It also adds another latency sample to the maximum the router waits for, pushing p99 further right. Only targeted queries convert extra shards into extra capacity.
  • The median latency is fine but p99 tripled after sharding. Why does the tail move first?
    A fan-out request completes only when its slowest participant answers, so every request draws N samples and takes the maximum. The median is set by typical shards, but the tail is set by whichever shard is momentarily busy — and with N shards you hit one of those far more often than before.
  • What does the SHARDING_FILTER stage in a shard's plan tell you?
    It is the stage that discards orphaned documents — copies a shard still holds after an interrupted or completed migration where cleanup has not run. It exists so a router-level read never returns the same document twice, and it is why reading a shard directly can show documents that the cluster-level query does not.
  • If a query genuinely cannot carry the shard key, what do you do?
    Contain it. Keep it off the latency-critical path, make sure a supporting index exists on every shard so each participant does a short index walk, bound the result size, and consider materialising the shape into a collection whose key makes the read targeted, refreshed by a background job.

saying these in an interview costs you the question

  • Proposes adding shards to fix slow broadcast queries
  • Reads only the top-level explain and ignores per-shard stats
  • Assumes fan-out latency is the average shard latency
  • Blames disk or CPU without checking shards touched per query
  • Thinks an index alone removes the fan-out

context