SysDesignPrep.com
Study guide 56 of 183

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

TechniqueHow it works
Replicate the keystore copies on several nodes (key#1 to key#N) and read a random one
Local in-process cacheeach app server caches the hottest keys for a second or two
CDN or edge cachingserve public content from the edge, far from the origin
Request coalescingmany concurrent misses for one key trigger one backend fetch
Read replicasroute 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.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.