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?
answer
- No shard key → broadcast to all N
- Latency = slowest shard, not average
- Fan-out QPS doesn't scale with more shards
- ORDER BY/LIMIT needs over-fetch; AVG from SUM+COUNT
- COUNT(DISTINCT) can't be summed
basics
~20 sThe 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 sWithout 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-- 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 mergedgo deeper
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.
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.
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.
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'