skip to content

In a horizontally sharded relational deployment, what happens when a query's filter does not contain the shard key, and how does its cost differ from a single-shard query?

level: middleimportance: must knowfreq 58%

answer

  1. No shard key → broadcast to all N
  2. Latency = slowest shard, not average
  3. Fan-out QPS doesn't scale with more shards
  4. ORDER BY/LIMIT needs over-fetch; AVG from SUM+COUNT
  5. COUNT(DISTINCT) can't be summed

basics

~20 s

The router cannot pick one shard, so it broadcasts to all of them and merges the results — a scatter-gather. Latency becomes the slowest shard's, work multiplies by shard count, and ORDER BY, LIMIT, and aggregates must be recombined by the router.

solid answer

~50 s

Without the shard key the router has nothing to prune on, so the query is broadcast to every shard and the results merged — scatter-gather (fan-out). Three costs matter: 1. **Latency is a tail, not a mean.** The request finishes when the *slowest* of N shards answers, so a rare slow shard becomes a common slow request. 2. **Throughput does not scale.** Each fan-out query consumes a connection and a worker on every shard, so adding shards does not add capacity for these queries — it adds load. 3. **Merge semantics.** `ORDER BY ... LIMIT k` needs each shard's top k and a router-side merge; `OFFSET` deep-paginates terribly; `AVG` must be rebuilt from SUM/COUNT; `COUNT(DISTINCT)` cannot be summed. A fan-out request also depends on *all* shards, so availability multiplies down. The fix is to make the second access path a point lookup — a directory/lookup table or a copy of the row keyed by the other dimension — or push it to a derived store.

code

text · 11 lines
text
-- client query
SELECT id, created_at FROM orders ORDER BY created_at DESC LIMIT 10;

router plan:
  shard_00 : SELECT id, created_at FROM orders ORDER BY created_at DESC LIMIT 10
  shard_01 : ... (same, LIMIT 10)
  ...
  shard_31 : ... (same, LIMIT 10)
  gather   : merge 320 rows on created_at DESC, emit first 10

with OFFSET 10000 LIMIT 10 each shard must return 10010 rows -> 320,320 rows merged

go deeper

for a junior

Know that the shard key decides which server answers, and that without it the query goes to every shard and the router merges the answers.

for a middle

Explain the three costs — tail latency, no throughput scaling, router-side merge — and name the merge traps: over-fetch for ORDER BY/LIMIT, AVG rebuilt from SUM and COUNT, COUNT(DISTINCT) not summable.

for a senior

Frame fan-out as a capacity and availability problem: a request touching all shards multiplies failure probability, so cap concurrency, set per-shard timeouts, and move the query to a lookup table or a CDC-fed derived store.

for a principal

Set a fan-out budget as an architectural constraint — which access paths are allowed to broadcast, at what QPS, and where the second access dimension lives — and make it a review gate rather than something discovered in an incident.

## The setup Horizontal sharding splits one logical table across N independent database servers. A routing layer — a proxy, a middleware tier, or logic in the application — inspects each query, extracts the **shard key**, and sends the query to the one shard that owns that key. That is what makes sharding scale: N shards, each seeing roughly 1/N of the traffic. That guarantee holds only when the router can determine the shard from the query. ## What happens without the shard key If the predicate contains no shard-key equality (or an `IN` list of keys), the router cannot prune. Its only correct move is to send the statement to **every** shard, collect the partial result sets, and combine them before returning. This is a **scatter-gather**, also called fan-out or a broadcast query. It is correct. It is just expensive in four separate ways. ## Cost 1 — latency becomes a tail statistic A fan-out request completes when the *last* shard replies. If any one shard is slow 1% of the time, a 32-way fan-out is slow about `1 - 0.99^32 ≈ 27%` of the time. You have converted a rare per-shard tail into a common per-request tail. Adding shards makes this worse, not better. ## Cost 2 — work amplification, so throughput stops scaling Each fan-out query occupies a connection, a worker thread, and buffer/CPU on all N shards. Single-shard traffic scales linearly with N; fan-out traffic does not scale at all, because every shard sees 100% of it. A workload that is 5% fan-out at 8 shards can be dominated by fan-out at 64 shards. ## Cost 3 — merge semantics are non-trivial The router must re-implement query operators above the shards: - `ORDER BY c LIMIT k` — each shard must return its own top k, and the router merges N×k rows to produce the true top k (over-fetch). - `OFFSET m LIMIT k` — each shard must return `m + k` rows so the router can skip correctly. Deep pagination gets brutal; use keyset ("seek") pagination on a sortable key instead. - `COUNT(*)`, `SUM`, `MIN`, `MAX` — sum or fold the partials. - `AVG` — cannot be averaged; the router must request SUM and COUNT and divide. - `COUNT(DISTINCT x)` — partials cannot be added at all; you need the union of values, or an approximate sketch such as HyperLogLog. - `GROUP BY ... HAVING` — needs two-phase aggregation: partial groups per shard, re-aggregate, then apply HAVING. - Unique/global ordering across shards has no common snapshot: each shard answers from its own transaction state at its own instant, so the merged result is not a point-in-time-consistent view unless the platform provides global snapshots. ## Cost 4 — availability couples to every shard If a request touches all N shards, its availability is roughly the product of theirs. One shard failing over degrades *every* fan-out query, so the blast radius of a single-shard incident becomes fleet-wide. Timeouts, bounded concurrency, hedged requests, and (where the product allows) partial results are damage control, not a fix. ## How to avoid it - **Keep the hot path keyed by the shard key.** Design the dominant access pattern around the routing dimension. - **Add a directory / lookup index.** Store a small mapping `alternate_id → shard_key` (often in its own store or sharded by the alternate id). The second access path becomes two point lookups instead of an N-way broadcast. - **Duplicate the row.** Maintain a second copy of the data keyed by the second dimension, written by one owner and propagated idempotently. - **Move the query out of the OLTP path.** Full-text search, reporting, and "find all X where Y" belong in a derived store fed by change data capture, not in a broadcast against the transactional shards. - **When fan-out is genuinely required** — admin tooling, low-QPS back-office, batch jobs — bound it: per-shard timeouts, a concurrency cap, iterate shard-by-shard for batch work rather than blasting all of them, and cache aggressively. ## When it is acceptable Low-QPS operations where the merged result set is small and each shard's slice is index-supported. The rule of thumb: fan-out is fine for *rare* queries and fatal for *hot* ones. Track a fan-out budget — what fraction of QPS touches more than one shard — and treat it as a capacity metric.

  • A user-facing endpoint must look up an order by its public reference code, but the orders table is sharded by customer_id. How do you serve it without fan-out?
    Maintain a lookup table mapping reference_code to customer_id (or directly to shard id), stored so that it can itself be found by reference_code — a separate small table sharded by the code, or a shared key-value store. The request becomes two point lookups instead of an N-way broadcast. Alternatively encode the shard or the customer id inside the reference code at generation time, so the router can derive the destination with no extra hop.
  • Why is OFFSET-based pagination so much worse under scatter-gather than on a single node?
    Each shard has no idea which of its rows fall inside the global offset window, so it must return offset+limit rows for the router to merge and skip. At offset 100000 across 32 shards the router pulls over three million rows to emit ten. Keyset pagination — carrying the last seen sort key and issuing WHERE (sort_key) < :last ORDER BY ... LIMIT k — keeps each shard's work bounded at k rows regardless of depth.

Asking a question with the shard key is calling the one colleague who owns the file. Asking without it is emailing all 32 colleagues and waiting for the last reply — you're blocked by whoever is at lunch.

saying these in an interview costs you the question

  • Saying fan-out latency is the average of the shards rather than the tail of the slowest one
  • Believing adding shards speeds up scatter-gather queries — it slows them and adds no capacity for them
  • Summing per-shard COUNT(DISTINCT) results, or averaging per-shard AVG results
  • Assuming the merged result is a consistent snapshot across shards
  • Treating a lookup/directory table as unnecessary because 'the router will figure it out'

context