A time-series table hashes writes on device_id, but one industrial sensor array alone generates 40% of all writes and keeps overloading its shard. What techniques let you avoid this single hot shard without abandoning hash sharding entirely?
answer
- hash of unchanging key always → same shard
- composite key = (entity, time_bucket)
- salting = random suffix, fan-in on read
- dedicated shard for known heavy tenant
- power-law traffic, not uniform
basics
~20 sChange the key so that one busy device's writes don't all pile onto one shard — for example, combine the device ID with a time bucket or a random suffix so its data spreads across several shards, then merge them back together when reading.
solid answer
~50 sThe core fix is making the shard key finer-grained so no single overloaded entity maps to exactly one physical shard. Common techniques: composite keys — hash on (device_id, time_bucket) so a busy device's writes spread across many shards over time instead of concentrating on one; key salting — append a random or round-robin suffix to the hot key before hashing, then fan-in the salted variants on read; dedicated/isolated shards — give a known heavy entity its own shard via a directory mapping instead of letting pure hashing decide; and write buffering/pre-aggregation — batch or roll up data before it crosses the shard boundary, reducing raw write volume from the hot entity. Every one of these trades read-side complexity (you now have to query multiple locations to reconstruct one entity's full data) for evenly distributed write load.
go deeper
Should grasp that one very active entity can overload its shard even in a hash-sharded system.
Should be able to name at least one concrete mitigation — composite key or salting — and describe roughly how it works.
Should articulate multiple mitigations, their read-side trade-offs, and how to detect a hotspot via per-shard metrics before it becomes an incident.
Should discuss proactively designing for power-law skew (dedicated tenant isolation, capacity planning per outlier) rather than only reacting after a hotspot appears.
## What a hotspot actually is A **hotspot** in a sharded system is a shard receiving disproportionate load relative to its siblings, and it's distinct from a **storage imbalance** — a shard can be perfectly average in total data size while still being overwhelmed on writes-per-second because one specific key on it is unusually active. Plain hash sharding maps a given key to exactly one shard deterministically: `hash(device_id) mod N` always returns the same shard number for that `device_id`, no matter how much traffic that device generates or how many total shards exist. If one device generates 40% of all writes, that 40% is permanently pinned to a single shard, and no amount of adding more shards to the cluster changes that, because the hash of that one unchanging key still resolves to one place. ## The techniques that actually move the load The fix has to change what value the hash function actually operates on, not how many shards exist. 1. **Composite key.** The most common technique is a composite key: instead of hashing `device_id` alone, hash on `(device_id, time_bucket)` — for example, `device_id` plus the current hour or day. Now the same device's data is spread across many distinct keys over time, each of which independently hashes to a (potentially different) shard, so the device's write volume gets distributed instead of concentrated. 2. **Key salting.** A second technique, key salting, is used when you can't naturally decompose the key by time or another dimension: you append a random or round-robin suffix (say, a digit 0–9) to the hot key before hashing, so writes for that one logical entity fan out across up to 10 shards. This works for any hot key, not just time-series ones, but it means reads for that entity must now query all salted variants and merge the results — you've traded a simple single-shard read for a scatter-gather read, in exchange for solving the write bottleneck. 3. **Directory-based isolation.** A third technique is directory-based isolation: rather than trying to force a known outlier entity through the general hashing scheme at all, you give it its own dedicated shard (or shards) via an explicit lookup mapping, sized and provisioned specifically for its load, while everything else continues through normal hashing. 4. **Reducing the raw write volume.** A fourth, complementary technique is reducing the raw write volume hitting the shard layer in the first place — buffering or pre-aggregating high-frequency writes (e.g., rolling 1-second sensor readings into 1-minute summaries) before they're persisted, so the shard sees fewer, larger writes instead of a flood of tiny ones. ## Why the problem exists at all The reason this class of problem exists at all is that real-world key distributions are almost never uniform — they tend to follow a **power law**, where a small number of entities (a viral post, a huge enterprise tenant, an industrial sensor array) generate vastly more activity than the median entity. A sharding scheme naively designed around "average" load per key will always be vulnerable to whichever entity turns out to be the extreme outlier, and outliers are usually impossible to predict at design time; they emerge from real usage. ## How it shows up in production Operationally, a hotspot shows up as one shard's latency, CPU, or throttling metrics diverging sharply from its siblings while the cluster-wide average looks fine — which is exactly why hotspots are dangerous: aggregate dashboards can hide them, and the team only notices once the hot shard starts rejecting or queuing writes. The failure mode compounds under **retry storms**: if writes to the hot shard start timing out, clients retrying those writes add even more load to the same already-struggling shard. ## Where it shows up by name This is a well-documented, named problem in managed sharded databases. - **AWS DynamoDB's** own documentation explicitly warns about "hot partitions" and recommends write sharding — appending a random suffix to a partition key for known high-traffic items — as the standard mitigation. - **Twitter's** engineering blog has historically discussed the analogous "celebrity problem," where a small number of accounts with enormous follower counts or posting volume require special-cased handling distinct from the sharding scheme that works fine for typical users. In both cases, the lesson is the same: hash sharding solves even distribution for a typical key distribution, but any system expecting a skewed, power-law distribution of activity needs an explicit strategy — salting, composite keys, or dedicated isolation — layered on top of the base hashing scheme to handle its outliers.
- What is 'key salting' concretely, and what does it cost on reads?It means appending a random or round-robin suffix — like a digit 0 through 9 — to a hot key before hashing, so writes for that one logical key get spread across several shards instead of one. The cost lands on reads: to reconstruct that entity's full data you now must query every salted variant and merge the results, trading a simple single-shard read for a scatter-gather read.
- Why doesn't simply adding more shards to the cluster fix a hotspot caused by one disproportionately busy key?Because a plain hash function still maps that one unchanging key to exactly one shard, deterministically, no matter how many total shards exist. Increasing the shard count redistributes the average load across more machines but doesn't change where that specific key's traffic lands — the fix has to change the key itself (via salting or a composite key), not the shard count.
It's like one extremely popular restaurant table that a single waiter is assigned to no matter how many extra waiters you hire — adding staff to the restaurant doesn't help until you actually split that one table's orders among several waiters (or give it a dedicated waiter of its own).
saying these in an interview costs you the question
- Suggests adding more shards alone fixes a single hot key
- Doesn't mention the read-side fan-in cost of salting
- Confuses a write hotspot (skewed load) with a storage imbalance (skewed size)
- Has no answer for how to detect a hotspot before it causes an outage