Hot keys, hot partitions and data skew
Why real traffic is never uniform, how a single hot key or partition takes down a sharded system, and the fixes: key splitting, local caching, request coalescing, write buffering, salting, adaptive rebalancing, and skew in joins and stream processing.
Reading is half of it. See this used in a real interview: walk through Design a Distributed Cache →
Sharding spreads load evenly only if keys are accessed evenly, and they never are. A celebrity's profile, a viral video, a flash-sale product, the live score of a final: real traffic follows a power law, where a tiny set of keys gets a huge share of requests. When one key lands on one shard, adding more shards does nothing. Interviewers love to ask "what happens when this key gets hot?", and having a toolbox ready is a strong signal.
Recognising skew
- Read-hot keys: many readers of one item (a celebrity profile, a trending post, a config value).
- Write-hot keys: many writers to one item (a like counter, a global sequence, a popular auction, a single partition key like "today's date").
- Partition skew: a poor partition key concentrates data, such as sharding orders by country or by timestamp.
- Processing skew: in joins and aggregations, one key carries most of the records, so one worker does most of the work.
Watch for it with per-shard and per-key metrics: request rate, CPU and latency per partition, and a top-K of the busiest keys. Averages across shards hide it completely.
Read-hot keys
| Technique | How it works |
|---|---|
| Replicate the key | store copies on several nodes (key#1 to key#N) and read a random one |
| Local in-process cache | each app server caches the hottest keys for a second or two |
| CDN or edge caching | serve public content from the edge, far from the origin |
| Request coalescing | many concurrent misses for one key trigger one backend fetch |
| Read replicas | route reads for hot shards to more replicas |
Even a one-second local cache turns a million requests per second on one key into roughly one request per second per app server. Request coalescing (also called single-flight) prevents a cache stampede when the hot key expires. See caching and Design a Distributed Cache.
Write-hot keys
Writes are harder, because copies must be combined:
- Split counters: write to one of N sub-counters (likes#0 to likes#15) chosen at random; read by summing them. Reads get slightly more expensive, writes scale N times. See counting at scale.
- Buffer and batch: aggregate increments in memory or a stream for a second, then apply one write.
- Queue the writes: put contended writes (bids, seat reservations) on a queue partitioned by key, processed by a single consumer in order, which turns contention into throughput. See inventory and flash sales.
- Shard the entity itself: split a huge inventory into buckets (100 tickets per bucket) so buyers spread across rows.
Bad partition keys
- Monotonic keys (timestamps, auto-increment ids) send all new writes to the last range partition. Fix with hashing, a prefix (a bucket number), or ids that do not sort by time at the top bits.
- Low-cardinality keys (country, status) produce a handful of giant partitions.
- Salting: append a random suffix to a hot partition key (date#0 to date#9) and read from all suffixes. This spreads writes at the cost of scatter-gather reads.
Choose keys with high cardinality and even access, and estimate the biggest single key, not just the total. See sharding and partitioning.
Adaptive approaches
- Split hot partitions automatically (DynamoDB and Bigtable split ranges that get hot).
- Detect hot keys at runtime with a top-K sketch and promote them to a replicated or locally cached tier. See probabilistic data structures.
- Isolate noisy tenants onto dedicated capacity. See multi-tenancy.
Skew in data processing
In Spark, Flink or a MapReduce job, one key with most of the records makes one task run for hours while the rest finish in minutes.
- Pre-aggregate with combiners before the shuffle.
- Salt the key for the first aggregation stage, then aggregate the partial results.
- Broadcast joins: send the small table to every worker instead of shuffling by a skewed key.
- Isolate the heavy keys and process them separately.
See batch and stream processing and Design an Ad Click Aggregator.
In the interview
After choosing a partition key, ask yourself out loud: "what is the hottest single key, and what happens to its shard?" Then pick the matching fix: local cache and replication for reads, split counters or queues for writes, and a better key for partition skew. Leaderboards (Design a Leaderboard) and ticket sales (Design Ticketmaster) are classic places to show this.
Checklist
- Per-key and per-partition metrics, with top-K detection.
- Read-hot: local cache, replicas, CDN, request coalescing.
- Write-hot: split counters, batching, single-writer queues.
- High-cardinality partition keys; no monotonic hot spots.
- Salting and broadcast joins for skewed processing.