Shard Keys — Hotspots, Fan-out Queries & Cardinality
The shard key is the product decision you cannot easily undo. Evaluate cardinality, skew, stability, and whether the hot transaction stays on one shard. Low-cardinality columns, monotonic time, and celebrity keys become hotspots. Queries that omit the key become fan-out.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Orders service key
Prefer
user_id, with whales isolated
The entity graph hangs off the user or the tenant. A single-order transaction and list-my-orders both stay local. Global open orders do not, so they leave the OLTP path.
- Rows that commit together share the key.
- A directory can map email to user id to shard.
- A mega-tenant moves to a dedicated shard instead of salting the whole product.
Alternative
created_at or order_id alone
Time piles inserts on the newest range. An order id makes one-row lookups easy and makes list-by-user a fan-out unless you add a secondary map.
- Write hotspot on the right edge if the key is time.
- Every user inbox or order list scatters if the key is the child id.
- Resharding a monotonic key does not fix the next minute of inserts.
Overview
Most sharding failures are key failures. A celebrity key, a low-cardinality column such as status or country, or a key that does not match the transaction (order lines on a different shard from the order) will show up as a hotspot or as a fan-out. The key is expensive to change. Treat the choice as a product decision.
This page is how you evaluate the key. Scatter-gather merge rules are the next page. A later change of key is a reshard, with dual-write and backfill, not an UPDATE of a column name.
By the end you should be able to:
- Score cardinality, skew, stability, affinity, and query alignment
- Name the four hotspot patterns and a mitigation that does not pretend salt is free
- Recognize a fan-out query that omitted the key
- Pick a key for orders and for chat without scattering the hot path
- Say which metrics prove the key is healthy
Decisions
- 1
Incoming query
- nextWHERE has shard key?
- ?
WHERE has shard key?
- YesRoute to 1 shard
- NoScatter-gather N shards
- 3
Route to 1 shard
- nextp50 and p99 stay local
- 4
p50 and p99 stay local
- 5
Scatter-gather N shards
- nextp99 is slowest plus merge
- 6
p99 is slowest plus merge
- nextDenormalize or go async
- 7
Denormalize or go async
Lesson map
Shard Keys — Hotspots, Fan-out Queries & Cardinality
The shard key is the product decision you cannot easily undo. Evaluate cardinality, skew, stability, and whether the hot transaction stays on one shard. Low-cardinality columns, monotonic time, and celebrity keys become hotspots. Queries that omit the key become fan-out.
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 q["Incoming query"] has["WHERE has shard key?"] single["Route to 1 shard"] ok["p50 and p99 stay local"] q -->|Incoming query to WHERE has shard key?| has has -->|Yes| single single -->|Route to 1 shard to p50 and p99 stay local| ok
Evaluation checklist
| Property | Good | Bad |
|---|---|---|
| Cardinality | Millions of distinct values | Tens or hundreds |
| Skew | Roughly even QPS and bytes | Power-law whales |
| Stability | Rarely changes | User changes org and the row must move |
| Affinity | Hot transactions stay on one shard | Every write touches two or more shards |
| Query alignment | Primary path includes the key | Primary path omits the key and fans out |
Cardinality is the number of distinct values. Too low and many rows collapse onto few shards. High cardinality with a whale is still skew. Measure both.
Hotspot patterns
- Monotonic keys. Time or an autoincrement as the only shard key. All inserts hit one range.
- Celebrity or whale tenant. One key owns a disproportionate share of traffic even when the hash is fair.
- Low cardinality. Sharding by
plan_typewith three values puts huge cohorts on three shards. - Launch spikes. A marketing campaign heats one geo list partition.
Mitigations: a salt or hash suffix on the hot key (reads of that key then fan out — say that cost); a dedicated shard for whales; adaptive splitting in the style of Vitess or TiDB; a cache and a queue in front of writes to the hot key. Salting is not the default for every key. It is a local escape hatch.
Fan-out queries
These omit the shard key:
SELECTorders where status is open, across users- Global search by email with no directory lookup
- Analytics
GROUP BYday on the OLTP shards
Avoid them with a secondary lookup (email to user id to shard), a denormalized query table, a CQRS read model, or by accepting that the work is batch or OLAP. The router should make scatter explicit. A silent scatter is how p99 dies.
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.
Imbalance in the Python score is the fattest bucket divided by the equal share. user_id on this toy set stays near uniform and clears the distinct-count bar. status and country do not. Entropy falls with them. The TypeScript planner is the production habit: if the filter lacks the key, return scatter and the shard count. Do not pretend it was a point lookup.
Orders: which key
| Key | Single-order transaction | List orders for user | Global open orders | Reshard ease |
|---|---|---|---|---|
| order_id | Great | Fan-out unless you keep a secondary | Fan-out | Medium |
| user_id | Great if the order stays with the user | Great | Fan-out | Medium |
| tenant_id | Good | Good per tenant | Fan-out and whale risk | Hard if whales share the hash |
| created_at | Hotspot on insert | Fan-out | Partial prune only | Poor |
Prefer user_id, or tenant_id with whale isolation, when the product graph hangs off that entity. Global open orders are not a reason to pick created_at. They are a reason to build a different read path.
Chat: a worked key
| Candidate | DM send | List inbox | Global moderation queue | Verdict |
|---|---|---|---|---|
| conversation_id | Single shard | Fan-out per user | Fan-out | Good if the inbox is a secondary index |
| user_id of the sender | Recipient copy is cross-shard | Easy for the sender | Fan-out | Usually a fan-out write |
| tenant_id | Fine for B2B | Fine | Whale risk | Fine plus dedicated whales |
Pattern that holds up: shard messages by conversation_id so a send is one shard. Maintain a per-user inbox projection for the list view. Do not scatter every time the app opens. The moderation queue is an async read model, not a query across every conversation shard on the request path.
Monitoring key health
- Rows, bytes, and QPS per shard, not only the global average.
- Top-K keys by write rate, from a sample.
- Cross-shard query ratio. The router should export a scatter ratio.
- Reshard readiness: a time-to-copy estimate for the largest shard.
If you cannot see per-shard QPS, you cannot see the whale until the shard is on fire.
Stripe's public DocDB talks are the same lesson at payments scale: the product engineer picks a shard key, the system splits that keyspace into chunks, and moving a hot chunk is an online operation. Search Stripe's engineering blog for the MongoDB and DocDB writeups. The InfoQ recording is a public walkthrough of that design. Uber's Schemaless post is the older lesson about cells and indexes around a shard key. Vitess documents how to select a vindex. TiDB's key-value guidance is the hotspot and split view.
Interview Q&A
What is shard-key cardinality?
Answer
The number of distinct values. Too low and many rows share a few shards. High cardinality is necessary and not sufficient. A single celebrity value can still own the traffic.
How do you find hotspots in production?
Answer
Per-shard QPS, CPU, and row counts. Top-N keys by sampled traffic. p99 broken out by shard. A healthy global average hides one shard at 100 percent.
Can you change a shard key later?
Answer
Yes, and it is painful. It is a dual-write, a backfill, and a cutover, which is the resharding page. It is not a metadata rename.
Why salt a hot key?
Answer
You spread one logical key across several physical buckets. Reads of that logical key must fan out and merge. Use it for a key you cannot give its own shard, and budget the read.
What is a directory or lookup table?
Answer
A map from a natural key, such as email, to the shard key, such as user id. Login does not scatter. The directory has to stay consistent with the move, or you 404 after a reshard.
What is the transaction-boundary rule?
Answer
Rows that must commit together live on the same shard key. Order and order lines. A message and its conversation row. If they must be atomic and they do not share a key, you have designed a saga.
Why is plan_type a bad shard key?
Answer
Three values means three shards no matter how many customers you have. Cardinality collapsed before you hashed.
user_id versus conversation_id for chat?
Answer
conversation_id keeps a send on one shard. The inbox is a projection per user. Sender-id sharding makes the recipient copy a cross-shard write on the hot path.
Pitfalls
- Hashing a column because it is on the table, not because the hot query filters on it.
- Calling salt a free performance win.
- Letting order lines and orders use different keys and then adding a cross-shard foreign key the engine cannot enforce.
- Watching only cluster-wide CPU.
- Shipping "search by email" as a scatter across user shards.
- Changing the key in a migration that only rewrites a column. The rows are on the wrong node.
An orders table has user_id, tenant_id, order_id, status, and created_at. One tenant is 40 percent of writes. The app lists a user's orders on every page view and runs a nightly open-order report. Pick the OLTP key, where the whale goes, and which query you refuse to scatter.
Go deeper
Vitess on selecting a sharding key is the MySQL-sharded reference. Uber's Schemaless post is a historical sharding lesson with secondary indexes around the cell. TiDB's key-value practices cover hotspot detection and splits. Stripe's DocDB design, in their engineering writing and in the InfoQ talk, is a chunk map over a chosen shard key.
Next: Cross-shard queries.