Cross-Shard Queries — Scatter-Gather, Aggregations & Avoidance
The day you shard, any query without the shard key becomes a scatter-gather. Latency is the slowest shard plus the merge, and every such query multiplies load by N. Merge sums and counts, never averages of averages. Keep scatter off the OLTP hot path.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Open orders on a user-sharded database
Prefer
A read model or a constrained fan-out
Denormalize, look up the shard, or project into a search or warehouse path. The checkout request never waits on every shard.
- Product lists and search sit on a projection.
- Ops tools can target a region list instead of all N shards.
- Analytics is a job, not a dashboard poll.
Alternative
Synchronous scatter on the request
Correct early, while N is small and an admin is the only caller. It becomes the capacity ceiling.
- One slow shard sets p99 for everyone.
- QPS is capped by the weakest shard once every query hits all of them.
- Each shard's snapshot is local. The merged sort is not one global snapshot.
Overview
The interview question is "how do you list all open orders in a user-sharded system?" The weak answer is SELECT everywhere, synchronously, on the request path. The strong answer is that the key and the query model were designed so the hot path never needs that, or the work is async or OLAP.
This page is merge semantics, aggregation traps, pagination, and avoidance. Choosing the key was the previous page. Moving rows between shards is the next one.
By the end you should be able to:
- Draw fan-out, partial results, and a merge
- Merge COUNT, SUM, AVG, DISTINCT, top-k, and percentiles correctly
- Rank sync scatter, cache, CQRS, and a warehouse
- Refuse offset pagination across moving shards
- Offer three designs for an open-orders API that are not a scatter poll
Flow
- 1
Query without shard key
- nextRouter fans out
- 2
Router fans out
- nextShard 0 subquery
- nextShard 1 subquery
- nextShard N subquery
- 3
Shard 0 subquery
- nextMerge, sort, limit, aggregate
- 4
Shard 1 subquery
- nextMerge, sort, limit, aggregate
- 5
Shard N subquery
- nextMerge, sort, limit, aggregate
- 6
Merge, sort, limit, aggregate
- nextLatency is slowest shard plus merge
- 7
Latency is slowest shard plus merge
Lesson map
Cross-Shard Queries — Scatter-Gather, Aggregations & Avoidance
The day you shard, any query without the shard key becomes a scatter-gather. Latency is the slowest shard plus the merge, and every such query multiplies load by N. Merge sums and counts, never averages of averages. Keep scatter off the OLTP hot path.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB c["Query without shard key"] r["Router fans out"] s0["Shard 0 subquery"] s1["Shard 1 subquery"] c -->|Query without shard key| r r -->|Router fans out to Shard 0 subquery| s0 r -->|Router fans out to Shard 1 subquery| s1
Cost model
- Latency is about the max shard latency plus merge CPU and network. One slow shard dominates p99.
- Capacity: if every query hits every shard, global QPS is about the minimum shard capacity. That is a death spiral as you add shards to "scale."
- Consistency: each shard returns its local snapshot. A global sorted merge is not a single snapshot unless you build one. Do not describe it as a single-node read.
Adding shards does not raise the ceiling of an all-shard query. It adds tails.
Aggregation pitfalls
| Aggregation | Naive scatter | Correct merge |
|---|---|---|
| COUNT of rows | Sum of counts | Sum is correct |
| SUM of x | Sum of sums | Sum is correct |
| AVG of x | Average of averages | Wrong. Merge sum and count |
| DISTINCT | Concatenate distinct sets | Distinct again on the merged set |
| ORDER BY with LIMIT k | Take one shard's page | Top-k from each shard, then merge |
| Percentiles | Average of percentiles | Wrong. Sketches such as t-digest, or merge raw samples |
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
The two shards with data have averages 20 and 80. Averaging those averages yields 50. The true average is 6000 / 150 = 40, because the second shard is smaller and hotter. An empty shard must not become a zero that you accidentally average in, and it must not divide by zero. Top-k is the same shape: each shard returns its local k, and the router keeps the global k. Shipping one shard's page is not a global page.
Avoidance, in the order you should prefer
- Denormalize, or a secondary table keyed by the query dimension.
- A lookup service: email to user id to shard.
- A CQRS read model. Projections in a search index or an OLAP store.
- Constrained fan-out. Only the shards a directory or a bloom says, not all N.
- Async reports. Never on the checkout path.
| Approach | Freshness | p99 | Ops | Prefer |
|---|---|---|---|---|
| Sync scatter-gather | Live | Poor at scale | Simple early | Rare admin tools |
| Cached aggregation | Seconds to minutes | Good | Invalidation | Dashboards |
| CQRS projection | Near-line | Good | Pipeline lag | Product search and lists |
| OLAP warehouse | Minutes to hours | Batch, not a request | ETL | Analytics |
Pagination
LIMIT 20 OFFSET 10000 across shards is expensive and unstable as rows move. Prefer keyset pagination with a global sort key and per-shard cursors, or paginate inside one shard only. Search systems learned this earlier. Elasticsearch documents why deep offsets are the wrong tool. The same arithmetic applies when the "index" is N databases.
Open orders, ranked
Bad: select every order with status open. That scatters N shards on every dashboard poll.
Better, in order:
- A status table that is still sharded by user id, so the status filter is local once you have the user. You still need user context. You did not solve the global list.
- A materialized open-orders index keyed by warehouse or region (a list partition) for ops tools. The fan-out is the regions you operate, not every user shard.
- An async export to a warehouse for analytics. The API returns 202 and a job id.
Hedging and partial results
For admin tools only: query shards under a budget and return a partial flag rather than blocking on the slowest shard forever. Never do that on a money-moving path. A partial total that looks complete is worse than an error.
Distributed SQL does not make joins free. Spanner and similar engines still pay network and coprocessor cost. Locality is still the design. A cross-shard join on the OLTP path is usually forbidden. Denormalize, or join in the app for a tiny fan-in you can name. "Select all" for an admin is a batch job with backoff, not the API default.
Interview Q&A
Why is an average of averages wrong?
Answer
Shards have different counts. Averaging their averages weights a small shard like a large one. Merge the sums and the counts, then divide once.
What limits scatter QPS?
Answer
Every query multiplies load by N. Global throughput collapses toward the weakest shard. Adding a shard adds a tail. It does not add capacity for an all-shard query.
How do search systems avoid OLTP scatter?
Answer
They keep a separate indexed read path, such as Elasticsearch, OpenSearch, or Meilisearch. The OLTP shards stay the source of writes. Search is a projection.
What about a cross-shard JOIN?
Answer
Usually forbidden on the OLTP path. Denormalize, or join in the application for a tiny fan-in. Distributed SQL still pays network cost for the same shape.
Are distributed SQL joins free?
Answer
No. They pay network and coprocessor cost. Design locality anyway. The engine hiding the scatter does not delete the scatter.
How should an admin select-all work?
Answer
A batch job with backoff and a cursor. Not the default of an API handler, and not a dashboard refresh timer.
Why does OFFSET across shards fail?
Answer
Each shard would skip a global offset it cannot see, or the router materializes a huge merge. Rows that move during a reshard change which page you are on. Use a keyset cursor or stay inside one shard.
When is a partial result acceptable?
Answer
On an admin tool with an explicit partial flag and a budget. Not on a path that moves money or shows a total a customer will treat as exact.
Pitfalls
- Averaging shard averages, or averaging percentiles, and publishing the number.
- DISTINCT that concatenates per-shard sets and counts duplicates twice.
- Believing a merged sort is a single-node snapshot.
- Deep OFFSET pagination as the public API.
- "We will scatter only until N is large." N being large is when scatter hurts.
- Hedging a checkout total so the slow shard is silently dropped.
A dashboard hits SELECT status = open every ten seconds against 32 user shards. Propose the three better options from this page, which one you would ship for warehouse ops, and which metric tells you the scatter ratio is climbing.
Go deeper
Vitess query rewriting is how a proxy turns a logical SQL statement into a scatter or a single-shard plan. Spanner's query execution overview is the distributed SQL cost model. Kleppmann chapter 6 covers partitioning and the queries that cross partitions. Elasticsearch's pagination guide is the keyset-versus-offset argument you can reuse across databases.
Next: Resharding and rebalancing.