A sharded relational cluster is running hot: one shard is saturated while the others have headroom. How do you determine whether the cause is data skew, request skew, or a single hot key, and what are your options for fixing each?
answer
- Bytes skew vs QPS skew vs one key
- Top-N keys at the router
- Hashing cannot split one key value
- Whale tenant gets its own shard
- Rebalance by load, not only by size
basics
~20 sMeasure three things per shard: bytes/rows stored, requests per second, and the top keys by request count. Even rows with uneven traffic means request skew; uneven rows means data skew; one key dominating means a hot key. Fix by moving buckets, isolating the heavy tenant, or splitting the key.
solid answer
~1 minDiagnose by separating the three signals, because they have different fixes: - **Data skew** — the shard stores far more rows/bytes. Usually a low-cardinality or badly distributed key, or a genuinely huge tenant. Fix by moving virtual buckets, splitting ranges, or giving the outlier its own shard. - **Request skew** — storage is balanced but QPS or CPU is not. One or a few keys generate a disproportionate share of traffic. Fix by isolating those keys (dedicated shard, sometimes called a whale shard), caching or replicating their read path, or adding read replicas for that shard only. - **Hot key** — a single shard-key value dominates. This is the hard one: hashing cannot help, because one value is atomic. You must either extend the shard key (`(tenant_id, bucket)` where bucket is a synthetic sub-key) so one logical entity spreads over several shards, or serve it from a cache/replica. Instrument per shard: rows, bytes, QPS, CPU, p99 latency, plus a top-N-keys counter sampled at the router. Without the top-N view you cannot tell request skew from a hot key. Prevention beats cure: model the expected distribution of the key before committing to it, and assume real-world tenant sizes follow a power law.
go deeper
Know that shards can become unbalanced and that the shard key's distribution is the usual cause.
Distinguish data skew from request skew and name the obvious remedies: rebalance buckets, isolate the big tenant.
Lead with the diagnostic method and the per-key instrumentation, then match a fix to each cause and state clearly that a hot key needs isolation or key extension, not rehashing.
Treat skew as inevitable and design for it up front: indirect placement so any key can be moved, heterogeneous shard sizing, quotas for abusive tenants, and load-aware rebalancing as an operational capability rather than an incident response.
## Why 'one shard is hot' is three different problems Sharding promises that N servers give you N times the capacity. It delivers that only if load divides evenly. In practice it rarely does, and the interview question is not "does skew happen" (it always does) but "can you tell the three kinds apart and treat them differently." **Data skew** is imbalance in what is *stored*: one shard holds a disproportionate share of rows or bytes. Symptoms are disk pressure, longer vacuum/compaction, a working set that no longer fits in that shard's buffer pool, and slowly degrading latency on that node only. **Request skew** is imbalance in what is *asked*: storage looks fine but one shard sees several times the queries per second, CPU or IO of its peers. Symptoms are a single node's CPU or connection pool saturating while disk usage is unremarkable. **Hot key** is request skew concentrated on one shard-key value — one tenant, one celebrity user, one partner integration hammering the API. It looks identical to request skew on a per-shard dashboard; only a per-key view distinguishes them, and the distinction matters enormously because the fixes are different. ## Diagnosing You need three layers of instrumentation, and most teams only have the first: 1. **Per shard**: rows, bytes on disk, QPS, CPU, IO wait, active connections, p99 latency. This tells you *that* something is unbalanced. 2. **Per key, top-N**: the routing layer already computes the shard key for every request, so it is the natural place to maintain a sampled top-N counter (a count-min sketch or a simple sampled histogram is enough). This tells you *whether one value* is responsible. 3. **Per key-and-operation**: is the heavy key heavy on reads or on writes? A read-heavy hot key can be absorbed by caching or a replica; a write-heavy one cannot. A quick decision rule: bytes skewed → data skew; bytes even and QPS skewed with a flat top-N → request skew spread over many keys on that shard (often a placement accident); QPS skewed with one key dominating the top-N → hot key. ## Fixing data skew If the key has adequate cardinality and placement is simply unlucky, rebalance: move virtual buckets from the heavy shard to lighter ones, or split the overloaded range. This is mechanical and online in mature systems. If the key has poor cardinality — someone sharded by `region` or `plan_type` — no rebalancing saves you; the ceiling is the number of distinct values. That is a shard-key change, i.e. a full data migration, and you should say so plainly rather than pretending a rebalance fixes it. If one tenant is legitimately enormous, accept heterogeneity: pin that tenant to its own shard with its own hardware sizing. Uniform shards are an aesthetic preference, not a requirement. ## Fixing request skew When many keys on one shard happen to be busy, bucket movement helps, but it must be driven by *load* rather than by *size*. A rebalancer that only equalises bytes will happily recreate the problem. If your platform only rebalances by size, you may need to move buckets manually using your own load metric. Read-heavy skew has extra levers that write-heavy skew does not: add replicas of the hot shard and route reads to them; add or fix caching in front of the hot access path; check whether an inefficient query (a missing index on that shard, an unbounded scan for the biggest tenant) is the real cause. Very often the "hot shard" is not a distribution problem at all but one bad query whose cost scales with a tenant's data volume. ## Fixing a hot key This is where candidates separate. Hashing cannot spread a single key value; by construction all its rows are on one shard. Your options: - **Isolate it.** Give the key its own shard with bigger hardware. Simple, effective, and admits that tenant sizes are power-law distributed. Directory-based routing makes this easy because placement is explicit per key. - **Extend the key.** Change the shard key from `tenant_id` to `(tenant_id, sub_bucket)` where the sub-bucket is derived from something the query knows — a sub-account, a device, a hash of the row id. The tenant now occupies K shards. This restores spread at the cost of turning that tenant's cross-bucket queries into fan-outs and breaking colocation-dependent transactions. It is a real schema and application change, not a config knob. - **Offload the read path.** Cache aggressively, or maintain a materialised, replicated copy of the hot entity's hot data. Works only for reads. - **Rate-limit or shape.** Sometimes the correct answer to one tenant consuming 60% of the cluster is a quota, not more hardware. ## Prevention Before choosing a key, run the distribution on real data: count rows and, if you can, requests per candidate key value, and look at the 99th percentile versus the median. If the top tenant is 100x the median, plan for isolation from day one. Keep placement indirect (virtual buckets or a directory) so you retain the ability to move a key without changing the hash function. And build the top-N key counter before you need it — reconstructing which tenant melted a shard from after-the-fact logs is miserable. ## How to answer Separate the three causes, name the metric that distinguishes them (per-key top-N at the router), then give the matching fix for each, and be explicit that a single hot key cannot be solved by hashing — only by isolation, key extension, or caching.
- Your largest tenant is 200x the median tenant. Would you shard around that, or handle it separately?Handle it separately. Designing the whole scheme around the outlier — for example splitting every tenant across sub-buckets — imposes cross-shard joins and lost colocation on the thousands of ordinary tenants who do not need it. The pragmatic answer is a directory or override that pins the whale to its own dedicated shard, sized for it, while the normal population stays hash-distributed with clean single-shard semantics.
- Why does a size-based rebalancer sometimes make a hot shard worse?It optimises the wrong objective. Moving buckets to equalise bytes can pull cold historical data off the hot shard while leaving the busy keys in place, or even move additional busy keys onto it because they are small. Rebalancing decisions need a load signal — requests, CPU time, or IO per bucket — not just stored size.
Evenly stocked shelves with one aisle permanently jammed: the fix depends on whether that aisle holds more goods, more shoppers, or one shopper buying everything.
saying these in an interview costs you the question
- Asserting that hash sharding prevents hotspots and therefore this cannot happen
- Not distinguishing storage skew from request skew, and prescribing one fix for both
- Claiming a single hot shard key can be spread by rehashing
- Assuming all shards must be identical in size and hardware
- Jumping to adding shards without first checking for a single inefficient query on that node