skip to content

A particular shard key value is receiving disproportionately more traffic than others - a 'hot' shard or hot key. What causes hotspots in a sharded system, and what techniques mitigate them?

level: seniorimportance: should knowfreq 65%

answer

  1. celebrity/power-user problem
  2. read hotspot means add cache
  3. write hotspot means key salting
  4. per-shard monitoring not cluster average
  5. DynamoDB adaptive capacity

basics

~20 s

A hotspot happens when way more requests hit one shard than the others, often because one key, like a celebrity's account or the newest data, is much more popular. Fixes include splitting that hot key's data across multiple shards or adding a caching layer in front of it.

solid answer

~50 s

Hotspots arise when the chosen shard key doesn't distribute load evenly in practice, even if it distributes data volume evenly - common causes are a celebrity/power-user key generating disproportionate traffic, sequential/time-based keys concentrating all new writes on one range, or a skewed access pattern where a small set of keys is read far more than others. Mitigations include key salting/splitting (append a random or rotating suffix to the hot key so its data and traffic spread across several shards, with reads fanning out and merging), caching hot reads in front of the shard so it only sees cache misses, dedicating extra capacity or a standalone shard to a known-hot tenant, and choosing a compound key up front that mixes a well-distributed component with the natural key. The right fix depends on whether the hotspot is a write hotspot (needs data redistribution) or a read hotspot (often solvable with caching/replicas alone).

go deeper

for a junior

Recognizes that one shard can get more traffic than others and that this is a problem.

for a middle

Can name the celebrity-key and sequential-key causes and suggest caching as a mitigation.

for a senior

Designs key-salting schemes, distinguishes read vs write hotspot fixes, and sets up per-shard monitoring.

for a principal

Builds or selects platform-level solutions such as adaptive capacity and tenant isolation/silo models, and sets policy for anticipating hot keys before launch rather than reacting after an incident.

## What a hotspot actually is A hotspot is what happens when the theoretical load-distribution property of your partitioning scheme fails to hold in practice: even though your hashing or ranging strategy assigns roughly equal shares of data or key space to each shard, the actual traffic - reads, writes, or both - concentrates heavily on one shard, turning it into a bottleneck while its siblings sit comparatively idle. It's a mismatch between how you partitioned the key space and how real-world access to that key space is distributed, and it's one of the most common ways a sharded system fails under real production load even after passing every load test that used synthetic, uniformly-distributed keys. ## Where hotspots come from Three causes show up repeatedly. 1. **The first is a skewed key popularity distribution**, often called the 'celebrity problem': in a social platform sharded by user ID, a celebrity account with millions of followers generates order-of-magnitude more reads and writes than a typical user, and no partitioning scheme distributes a single key's traffic across multiple shards by default - the celebrity's data lives on exactly one shard, so that shard absorbs the celebrity's entire load regardless of how evenly other users' keys are spread. 2. **The second is temporal/sequential clustering:** if the shard key is or correlates with a timestamp or auto-incrementing ID, all new writes land on whichever shard owns the current top of the range, making 'newest data' a permanent write hotspot even though total data volume across shards stays balanced over time. 3. **The third is workload skew unrelated to the key itself** - a marketing campaign, a viral post, or a retry storm from a misbehaving client can suddenly concentrate load on one or a few keys that were previously unremarkable. ## Read hotspots versus write hotspots The mitigation strategy depends heavily on whether the hotspot is a read hotspot or a write hotspot, because they have very different fixes. - **Read hotspots are the more tractable case.** Because reads don't have to be linearizable with every other read, you can absorb them with a cache in front of the shard (Redis, Memcached, or a CDN for read-mostly public data) so only cache misses ever reach the database, or by adding more read replicas specifically for that shard and routing reads to whichever replica is least loaded. - **Write hotspots are harder**, because every write ultimately has to land somewhere durable, and you can't cache your way out of write volume. The standard technique here is **key salting** (sharding key splitting): instead of writing all of a hot key's data under one literal key, append a random or rotating suffix (e.g. `celebrity123#0` through `celebrity123#9`) so the data, and critically the write traffic, spreads across N shards instead of one. Reads for that logical key then have to fan out across all N salted variants and merge results, trading some read complexity for write scalability; this pattern is explicitly recommended in DynamoDB's and Bigtable's own scaling guidance for exactly this reason. ## Dedicating capacity to a known-hot key Another mitigation is dedicating extra or isolated capacity to a known-hot key up front rather than reactively - if you can identify in advance that a specific tenant or entity will be disproportionately large (a known enterprise customer signing a big contract, a flagship product SKU), placing it on its own dedicated shard or shard group avoids it ever contending with the noisy neighbors sharing a generic hash bucket; this is a common pattern in multi-tenant SaaS, sometimes called **tenant isolation** or a 'silo' model for the largest tenants while smaller tenants share 'pooled' shards. ## The structural fix: a compound shard key A structural mitigation that prevents many hotspots before they occur is choosing a compound shard key that mixes a naturally meaningful component with a well-distributed one from the start - e.g. instead of sharding purely by `customer_id` (which lets one giant customer dominate a shard) or purely by a random hash (which loses any locality), some systems shard by a combination like `hash(customer_id) mod small_bucket_count` combined with the ID itself, deliberately spreading even a single very large customer's data across a bounded number of shards rather than one. ## How it looks in production In production, hotspots show up as one node in a monitoring dashboard with dramatically higher CPU, IOPS, or queue depth than its siblings while aggregate cluster metrics look fine - this asymmetry is exactly why per-shard, not just cluster-average, monitoring and alerting is essential; averaging across shards can completely hide a hotspot that's actively degrading service for the users unlucky enough to be routed to it. Bigtable's public documentation and DynamoDB's adaptive capacity feature, which automatically shifts extra throughput capacity toward a hot partition on the fly, both exist specifically because this failure mode is common enough to warrant first-class platform support rather than leaving every team to invent salting logic independently.

  • Why doesn't adding more shards to the cluster automatically fix a hotspot caused by one extremely popular key?
    Because that single key's data still lives on exactly one shard under a standard hash or range partitioning scheme. Adding shards redistributes other keys but does nothing for the one shard that already owns the hot key, so the hotspot persists until you specifically split or salt that key's data across multiple shards.
  • What's the trade-off introduced by key salting to fix a write hotspot?
    Salting spreads writes for a logical key across multiple physical shards by appending a random or rotating suffix, which fixes the write imbalance, but it means a read for that logical key now has to fan out across every salted variant and merge the results, trading write scalability for extra read complexity and latency.
  • Why might averaging load metrics across all shards fail to reveal a serious hotspot problem?
    Because a cluster-wide average blends one severely overloaded shard with many lightly loaded ones, producing a healthy-looking average even while users hitting the hot shard experience real degradation. Per-shard, or per-partition, monitoring is needed to catch the asymmetry.

A hotspot is like one toll booth on a ten-lane highway getting all the traffic because a popular exit is right after it, while the other nine lanes stay empty - widening the whole highway doesn't help until you specifically add capacity or reroute traffic around that one booth.

saying these in an interview costs you the question

  • proposes 'just add more shards' as the fix for a single hot key
  • doesn't distinguish read hotspot mitigation from write hotspot mitigation
  • suggests salting without mentioning the read fan-out cost it introduces
  • only monitors cluster-average metrics
  • assumes hash partitioning makes hotspots impossible

context