Hiring BarSupport

Hot Keys

Pattern50 min read8 diagrams

Technologies that implement this pattern: Redis · DynamoDB · Cassandra · Apache Kafka · Flink · PostgreSQL

Why This Matters#

Every partitioning scheme has a key it can't split. Hashing, ranges, consistent-hash rings, directories: they all spread keys across nodes, and none of them spreads load on one key. The celebrity profile, the Super Bowl ad, the flash-sale SKU, the tenant that is 40% of traffic and the counter for today's date each land on exactly one partition, one Redis thread, one Kafka partition, one Flink subtask. The cluster can be 2% utilized while that one node is on fire and every other key it hosts times out with it.

Most candidates treat hot keys as a sharding question: "add more shards." Staff engineers treat them as a skew-budget question. The design question is not how many shards but what is the hottest single key this system will ever see, what is one partition's ceiling, and which fix do I reach for when the first exceeds the second? Real traffic is Zipfian. Plan for the top 0.01% of keys explicitly, because no hash function will do it for you.

This page is the index. The library mentions hot keys in about 40 pages, and each page fixes its own key: the celebrity cache key in Distributed Cache, the viral counter in Leaderboard, the celebrity fan-out in News Feed, the hot partition in DynamoDB. What this page adds is the part those pages assume: how you find the key, how you tell read-hot from write-hot, which fix to use in what order and at what price, and when the right fix is a product decision rather than an infrastructure one.

A Staff answer names the hot key before the interviewer does. If you draw a sharded store and the interviewer has to say "what about Taylor Swift?", you've already lost the signal.

The 60-Second Version#

  • Name the hottest key and its peak rate before you finish the high-level design. Compare it to one partition's ceiling: ~100–200K simple ops/sec per Redis shard, 1,000 writes/sec and 3,000 reads/sec per DynamoDB partition, ~10 MB/s per Kafka partition (conservative), one CPU core per Flink subtask. If the hottest key is above ~50% of that ceiling at peak, you have a hot-key design problem. Below it, you don't.
  • Classify it: read-hot, write-hot, or invariant-hot. Read-hot keys are solved with copies (L1 cache, replicas). Write-hot keys with commutative updates are solved by splitting (salting, pre-aggregation). Write-hot keys that guard an invariant (inventory, a balance) are a Concurrency Control problem, and the fixes change completely.
  • Fixes have a cost order. Climb it, don't jump. Cheapest first: request coalescing (free, no staleness) → L1 in-process cache (1–5s staleness) → key replication across N shards (N× write and invalidation fan-out) → write salting / pre-aggregation (N× read fan-in, approximate totals) → splitting the entity (schema and product change) → a dedicated cell (infrastructure and on-call).
  • Detect with sampling and top-K, not with logs. Sample 0.1–1% of requests per host into a count-min sketch or space-saving top-K, report every 1–10 seconds, and alert on shard.cpu skew > 2–3× the median. Without per-key visibility, a hot key looks like "shard 14 is slow."
  • The blast radius is the neighbors. A hot key at 38K writes/sec doesn't just slow itself; it takes the p99 of a million co-located keys from 0.3 ms to 40 ms. Hot-key protection is a platform concern because the victims belong to other teams.
  • Celebrity accounts are a product decision. Whether a 100M-follower account gets fan-out-on-read, a 2-second-stale like count, or a dedicated cell is a tier policy that product signs off on, and it's set before the celebrity arrives, not during the incident.

The Problem#

Real access is skewed. A 50M-user social app has a few hundred accounts with more than 10M followers; a marketplace has 0.1% of SKUs taking 30% of views during a sale; a SaaS platform has one tenant that is a third of all traffic. Under a hash partitioner each of those keys lives on one partition with a fixed ceiling. When a live event drives one cache key from 2K to 450K reads/sec at 3 KB per value, the shard's 10 Gbps NIC saturates at roughly 400K reads/sec, and every other key on that shard starts timing out. When a viral post takes 38K likes/sec into one Redis INCR, the single command thread pins at 95% CPU. Adding shards does nothing for either, because the key still hashes to one place. The job is to know which keys will do this, detect the ones you didn't predict within seconds, and have a ranked set of fixes ready, each with a known cost and a named owner.


Case Studies That Use This Pattern#

Every hot key in the library, indexed by what kind of hot it is. Go to the linked page for the system-specific fix.

PageThe Hot KeyKindThe Fix That Page Uses
Distributed CacheLive-event summary key, 450K reads/s on one shardRead-hotAutomatic near-cache, N-suffix replication
Scaling ReadsCelebrity profile at 1M reads/sRead-hotSampled top-K → L1 with 2s TTL
Content Delivery NetworkViral object on one origin shieldRead-hotTiered cache, request collapsing at the edge
URL ShortenerA link shared in a Super Bowl adRead-hotEdge cache + L1, immutable mapping
Search TypeaheadSingle-letter prefixesRead-hotPrecomputed top-K, client and edge cache
LeaderboardViral post's like counterWrite-hot, commutativePre-aggregation, sharded counter
Ad Click AggregationOne ad getting 40% of clicksWrite-hot, commutativeTwo-stage aggregation, salted keys
Scaling WritesHot item in DynamoDB/CassandraWrite-hot, commutativeWrite salting with discoverable bucket count
News FeedCelebrity post fan-outWrite-amplification-hotPull for the heavy tail (hybrid fan-out)
Chat AppA 100K-member groupFan-out-hotShared group log, per-member cursors
Ticket DropsOne SKU's inventory rowInvariant-hotAdmission control, bucketed inventory
Concurrency ControlA hot coupon or seat mapInvariant-hotSplit state, single-writer, admission in front
Reservation SystemsA popular event's seat mapInvariant-hotHolds with expiry, queue in front
Rate LimiterOne tenant's counter at 40% of trafficWrite-hot, coordinationLeased budgets, sub-buckets
Apache KafkaSkewed producer keyOrdering-hotKey redesign, sub-keys where order allows
Sharded DatabaseWhale tenantTenant-hotDirectory move to a dedicated shard

The table is the first teaching point: "hot key" is at least five different problems. The fixes for a read-hot profile and an invariant-hot inventory row share nothing except the symptom.


The Four Kinds of Hot#

The phrase "we have a hot key" hides incompatible problems. Classify before choosing a fix.

KindExampleConstraintStrategyFailure ModeCorrectness Bar
Read-hotCelebrity profile, live-event summary, config blobOne shard's CPU or NICCopies: coalescing, L1, key replication, CDNStale copies; invalidation fan-outStale ≤ 1–5s, product-signed
Write-hot, commutativeLike counter, view count, click totalsOne partition's write ceilingSplit: salting, pre-aggregation, local batchingApproximate or delayed totals; N× read fan-inEventually exact (seconds)
Invariant-hotLast 200 units of a SKU, an account balance, a seatOne row must enforce a ruleAdmission in front, bucketed state, single writerOversell, or "sold out" with stock leftExact, never violated
Tenant-hot / ordering-hotWhale customer, one Kafka key, one Flink keyAll of one entity's work on one nodeIsolation: dedicated shard/cell, sub-keys where order allowsNoisy neighbor; lag on one partitionPer-entity ordering preserved

🎯 Staff Move: "Before I pick a fix I want to know what kind of hot this is. The celebrity profile is read-hot, and I'll fix it with copies. The like counter is write-hot but commutative, so I can split it. The inventory row is invariant-hot, and splitting it changes the oversell math, so it goes through admission control. Three hot keys, three different fixes."

Read-Hot vs Write-Hot at a Glance#

Diagram: Read-Hot vs Write-Hot at a Glance

The asymmetry is the point. Read-hot fixes multiply copies and pay in staleness. Write-hot fixes multiply keys (or collapse writes) and pay in read fan-in and exactness. A candidate who proposes "replicate the counter" for a write-hot key has made every write N writes, which is the opposite of what's needed.


The Core Tradeoff#

The ranked menu. Ordered by total cost: build effort, runtime cost, and how hard it is to undo. Climb it from the top.

RankFixWhat WorksWhat BreaksWho Pays
0Do nothing (headroom)Hottest key < 50% of partition ceiling at peakA celebrity event you didn't forecastThe on-call, at 3am, if the forecast was wrong
1Request coalescing (singleflight, leases)N concurrent misses → 1 origin call; no staleness addedDoesn't reduce cache-tier load, only origin load; per-host scopeNobody much: a few ms of wait for joiners
2L1 in-process cache (hot keys only)Removes the network hop; 1M reads/s → hosts ÷ TTL on the shardUp to TTL of cross-host disagreement; memory per hostUsers who see values flip between hosts for 1–5s
3Key replication (k#0..k#N-1, read any)Read ceiling × N; survives a shard lossEvery write and every invalidation × N; partial-update windowsWrite path latency; the team debugging "replica 3 is stale"
4Write salting / sharded counterWrite ceiling × NReads fan in N sub-keys; totals are sums; N is hard to lower laterRead latency; product (counts are seconds behind)
5Pre-aggregation / local batching38K writes/s → 1 write/s per host or per windowLoses per-event granularity at the store; crash loses ≤1 windowProduct (approximate live counts); stream team owns the job
6Split the entitySummary vs detail, per-region shards of one board, per-bucket inventorySchema change; product semantics change ("top 100 per region")Product and API consumers; migration owner
7Dedicated cell / tierFull isolation: the hot entity can't hurt neighborsSeparate infra, deploys, capacity planning, on-callPlatform headcount and $/month; the routing layer's complexity

Two things to say about the ladder in an interview. First, ranks 1 and 2 should be on by default in the platform client, because they cost almost nothing and they're the difference between a hot key and an incident. Second, ranks 4 and up are one-way-ish doors: once readers sum 16 sub-keys, lowering N is a migration, and once product ships "top 100 per region," going back to a global board is a product change.

Diagram: The Core Tradeoff

Staff Default Position#

Name the hottest key and its peak rate up front; ship coalescing and automatic hot-key L1 as platform defaults; climb the ladder per key only when measured load passes 50% of one partition's ceiling.

The default stack: every client of a shared store coalesces concurrent misses and samples keys into a top-K sketch. Keys above a fleet-wide threshold (for example, 0.5% of a host's reads, or an estimated 20K reads/sec fleet-wide) are promoted automatically into an in-process L1 with a 1–5 second TTL. Write-hot commutative keys are salted or pre-aggregated, with the bucket count stored in metadata readers can discover. Invariant-hot keys get admission control in front rather than splitting. Celebrity and whale entities are identified by a tier attribute that product owns, and their behavior (pull instead of push, stale counts, a dedicated cell) is decided in advance. Every promotion is visible on a dashboard and every manual override has an expiry.


When to Deviate#

  • Values that can't be stale even for a second. Prices at checkout, permissions after revoke, a balance. An L1 copy is a correctness bug. Serve from the primary with coalescing only, or use versioned keys with synchronous invalidation, and add capacity to that partition (bigger node, read replicas of that shard).
  • Predictable events. A product launch, a scheduled final, a ticket on-sale at 10:00. Don't wait for detection: pre-mark the keys hot from the event calendar and pre-warm L1 and replicas 15 minutes before. Detection exists for the events you didn't predict.
  • Small fleets. With 8 app hosts, an L1 cache only cuts a hot key to 8 ÷ TTL reads/sec at the shard, but each host is also 12% of traffic, so a 2s L1 TTL buys less. Prefer key replication or a bigger shard.
  • Ordering-hot keys. If one Kafka key or one Flink key must stay ordered, salting it breaks the ordering guarantee. Either prove the order only matters within a sub-entity (per-user-per-chat rather than per-chat) or accept the single partition and give it a bigger consumer.

One Question, Three Levels#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First moveShards by a high-cardinality key and moves onNames the hottest key, its peak rate, and the partition ceiling before finishing the designAsks which teams will hit the same key class (celebrities, whales, events) and whether the tier policy should be org-wide
Detection"We'd see it in the metrics"Per-key sampling into top-K, shard skew alerts, a dashboard of promoted keysHot-key telemetry is a platform contract: every shared-store client emits it, one dashboard for the org
Fix choice"Add a cache" or "add shards"Classifies read/write/invariant-hot and climbs the cost ladder; names staleness and fan-in costsMakes ranks 1–2 automatic defaults and reserves ranks 6–7 for a reviewed exception process
Celebrity handlingSame code path for everyoneThreshold-based hybrid (push below 1M followers, pull above)Celebrity tier is a product-owned policy with a review board, hysteresis, and a cost line
Failure"Retry on timeout"Knows retries double load on the hot shard; caps retries and isolates neighborsFunds a hot-event game day each quarter; uses the event calendar to pre-mark keys
OwnershipWhoever owns the storeService team owns key design; platform owns detection and L1Redraws the contract: platform owns automatic mitigation, product owns tier policy, events team owns the calendar
Why "First move" separates levels

The L5 answer, "shard by user_id," is correct and still gets downleveled, because it says nothing about the one user_id that will be 1,000× the median. The Staff candidate puts a number on it: "the hottest account has 80M followers; a post drives about 300K reads/sec for the first minute, and a Redis shard does about 150K, so I need copies." The Principal candidate notices that the same celebrity breaks the feed, the counter, the notification fan-out and the search index on the same day, and asks for one tier policy instead of four separate fixes.

Why "Fix choice" separates levels

"Add shards" fails because one key lives on one shard. "Add a cache" works for read-hot keys and fails silently for write-hot ones: caching an INCR doesn't reduce writes. The Staff answer classifies first, then names the price of each fix: L1 costs seconds of staleness, salting costs N-way read fan-in, splitting costs a product change. The Principal answer recognizes that the cheap fixes should never be a per-team decision. They belong in the client library so the 40th team doesn't rediscover them during an incident.

Why "Celebrity handling" separates levels

A threshold in code ("pull if followers > 1M") is a Staff answer and a good one. But the threshold decides latency, cost and feature parity for the most visible users on the platform, and it flaps when an account crosses it during a viral moment. The Principal answer makes the tier explicit, owned by product, with hysteresis (enter at 1M, leave at 800K), a documented feature gap ("celebrity posts may appear in feeds up to 30s later"), and a cost line that shows what the tier saves.


Where the Design Splits#

#Fault LineThe Tension
1Detect vs Pre-declareReact to measured heat (works for surprises, lags by seconds) or mark keys hot in advance (instant, needs a calendar and an owner)
2Generic Fix vs Entity-Specific FixOne platform mechanism for every hot key, or a fix shaped to the entity's semantics
3Absorb vs RefuseSpend infrastructure to serve all the heat, or cap it with admission, shedding or product limits
4Automatic vs Human PromotionLet the system promote and demote hot keys, or require an operator

Fault Line 1: Detect vs Pre-declare#

Detection catches the keys nobody predicted, but it lags. A 10-second sampling window plus a config push is 15–40 seconds in which the shard is already saturated, and the Distributed Cache celebrity-key incident shows what 3 minutes of lag costs. Pre-declaring is instant but needs an owner: someone has to know the final starts at 20:00 and the SKU drops at 10:00. Who pays: detection-only pays in the first 30 seconds of every surprise event; pre-declaring pays in a process that product and events teams must keep current. Staff default: both. Pre-mark known events from a calendar, detect everything else, and treat any incident caused by a known event as a process miss. Deviate when: traffic has no scheduled events (B2B APIs), and detection alone is fine.

Fault Line 2: Generic Fix vs Entity-Specific Fix#

A generic L1 promotion works for any read-hot key without the owning team doing anything. But some heat has structure the platform can't see: a leaderboard can be split by region, a summary can be separated from its detail, a group chat can switch from per-member inboxes to a shared log. Those entity-specific fixes are cheaper at runtime and more permanent, and they need the domain team. Who pays: generic fixes are paid for by the platform team and by users who see 1–5s of staleness; entity-specific fixes are paid for by the domain team's roadmap. Staff default: generic for the first response, entity-specific when the same entity is hot every week. Deviate when: the entity is invariant-hot, where no generic copy-based fix is safe.

Fault Line 3: Absorb vs Refuse#

You can buy capacity for the full heat (replicas, a dedicated cell, a bigger node) or you can refuse some of it: admission control in front of a sale (Ticket Drops), rate limits per tenant (Rate Limiter), or a product limit like "live like counts update every 2 seconds." When 600K people want 20K seats, absorbing all 600K attempts at the inventory row is the mistake. Who pays: absorbing pays in $/month and idle capacity between events; refusing pays in user-visible waiting, throttling, or coarser features. Staff default: absorb read heat with copies (cheap), refuse invariant heat with admission (correct), and negotiate the product limit for write heat on vanity metrics. Deviate when: the hot entity is a paying whale whose contract specifies throughput; then you absorb it in a dedicated cell and bill for it.

Fault Line 4: Automatic vs Human Promotion#

Automatic promotion reacts in seconds and never sleeps. It can also flap: a key crosses the threshold, gets promoted, its shard load drops, it gets demoted, and the load comes back. Human promotion is deliberate and too slow at 3am. Who pays: automatic flapping pays in latency oscillation and confusing dashboards; manual pays in incident minutes. Staff default: automatic promotion with hysteresis (promote at 20K reads/s, demote below 5K reads/s sustained for 5 minutes), a manual override with an expiry, and an audit log of every promotion. Deviate when: the fix is expensive or semantic (salting, splitting, moving to a cell); those stay human decisions.

Diagram: Fault Line 4: Automatic vs Human Promotion

The lifecycle is the talking point: promotion has entry and exit conditions with different thresholds, so a key at 6K reads/sec doesn't oscillate between states every few seconds.


Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"We'll shard it, so load is spread evenly""Doesn't know one key lives on one shard""Hashing spreads keys, not load on a key. The hottest account is about 300K reads/sec and one shard does about 150K, so that key needs copies."
"Consistent hashing handles hot spots""Confuses placement with load""Consistent hashing helps when nodes join and leave. It does nothing for a single key. That's a replication or caching problem."
"Put a cache in front of the counter""Hasn't separated read-hot from write-hot""Caching the read side helps; the 38K INCR/sec still hits one thread. I'll batch locally and flush once a second, or salt into 16 sub-keys."
"Split the inventory row into 32 buckets""Hasn't priced the invariant""Bucketing raises the ceiling 32× but 'sold out' becomes a sum, and the last units strand in empty buckets. For a 10× oversubscribed drop I'd put admission control in front first."
"Auto-scaling will handle the spike""Thinks capacity is fungible""Auto-scaling adds nodes for new keys. The hot key's partition doesn't get bigger. DynamoDB's own paper says a single hot item won't benefit from a split."
"We'll treat celebrities the same as everyone""Hasn't considered product tiers""Above about 1M followers I switch to pull on read, with hysteresis so accounts don't flap. Product signs off that celebrity posts can arrive a few seconds later."
"We'd find the hot key from the logs""Detection after the fact""Every client samples 1% of keys into a top-K sketch and reports every 10 seconds. The alert is shard CPU skew over 3× median, and the dashboard names the key."

Quick Reference#

The decision tree. Start from the measured or forecast peak of the single hottest key.

Diagram: Quick Reference

Staff Sentence Templates#

"The hottest key in this design is [entity], at about [N] [reads/writes] per second at peak. One [Redis shard / DynamoDB partition / Kafka partition] handles about [M], so that key is [N/M]× over a single partition and I need [copies / splitting / admission]."

"This key is [read-hot / write-hot / invariant-hot], so the fix is [L1 + coalescing / salting into N sub-keys / admission control]. The price is [1–5s of staleness / N-way read fan-in / users waiting in a queue], and [product / the service team] has signed off on that."

"Coalescing and hot-key L1 are on by default in the client, so the first [30] seconds of a surprise event are absorbed without a human. Anything past rank [3], salting, splitting or a cell, is a deliberate change with a named owner."

"Celebrity accounts are a product tier, not a code path. Above [threshold], [behavior]; product owns the threshold and the feature gap, and the platform owns making the tier cheap."


Detection: Finding the Key Before It Finds You#

You cannot fix a hot key you can't name. Detection has three layers, and each catches what the others miss.

LayerMechanismLatency to Name the KeyCostMisses
Server-side skewPer-shard CPU, ops/sec, NIC bytes; alert on skew > 2–3× median30–60s (metric scrape)FreeWhich key: tells you shard 14, not event:9921
Server-side key samplingredis-cli --hotkeys (requires an LFU eviction policy), store-provided contributor insights, MONITOR samplingMinutes, often manualLow, but some tools are heavy on a busy nodeNot continuous; operator must think to run it
Client-side sampled top-KEach host samples 0.1–1% of keys into a count-min sketch or space-saving counter, reports top 50 every 1–10s1–10s, automatic~100 KB memory per host, negligible CPUKeys that are hot across many hosts but small on each, unless aggregated centrally

The client-side layer is the one that enables automation, because it names the key fast enough for the client to act on it. A count-min sketch gives a frequency estimate in fixed memory: d hash rows of width w, an error of about ε × total where w = e/ε, with probability 1 − δ where d = ln(1/δ). With w = 2,048 and d = 4, the sketch is 32 KB of 32-bit counters and overestimates any key's count by at most ~0.13% of the window total with ~98% confidence. That is enough precision to find any key above 0.5% of traffic, and that is the only question detection has to answer.

Diagram: Detection: Finding the Key Before It Finds You

Two detection rules that matter more than the algorithm. First, normalize by bytes as well as by count: a 3 KB key at 400K reads/sec saturates a 10 Gbps NIC while the CPU looks fine, so report bytes/sec per key alongside ops/sec. Second, aggregate centrally. A key at 0.2% of every host's traffic is invisible locally and is 100% of one shard's load when 500 hosts agree on it.

🎯 Staff Insight: Detection exists for the surprises. For anything on a calendar, such as a launch, a final, a ticket on-sale or a scheduled push notification, the keys are known in advance. Pre-mark them hot 15 minutes early and treat a hot-key incident on a scheduled event as a process failure, not bad luck.


Implementation Deep Dive#

Each implementation is the cross-cutting version. The linked pages show the system-specific variant.

1. Request Coalescing — Singleflight per Host, Leases per Key#

The cheapest fix and the one that should never be optional. When 10,000 requests miss on the same key at the same moment, only one per host (singleflight) or one per key cluster-wide (lease) goes to origin.

inflight = {}                                  # per process: key -> future

function get(key):
    v = cache.get(key)                         # may return value, or LEASE token, or WAIT
    if v.is_value: return v.value
    if v.is_wait:                              # another client holds the lease
        sleep(5ms + jitter); return get(key)   # usually filled within a few ms
    if key in inflight: return inflight[key].await()
    fut = inflight[key] = new Future()
    row = origin_limiter.run(() => db.read(key))
    cache.set_with_lease(key, row, v.lease_token, ttl=jitter(60s))
    fut.resolve(row); delete inflight[key]
    return row

Numbers: per-host singleflight turns a 10,000-request miss storm across 500 hosts into ≤500 origin calls; a server-side lease that hands out one token per key per few seconds turns it into ~1. Facebook's memcache paper reports exactly this kind of collapse from leases (see In the Wild). Coalescing does not reduce the read load on the cache shard itself; that needs rank 2. Detail: Read-Heavy Systems, Caching Fundamentals.

2. Automatic Hot-Key L1 — Promote on Detection, Demote with Hysteresis#

PROMOTE_QPS = 20_000      # fleet-wide estimate
DEMOTE_QPS  =  5_000      # sustained for 5 min
L1_TTL      = 2s          # product-approved for this namespace

on promote_list_update(keys):           # pushed by hot-key service
    for k in keys: l1.allow(k, ttl=L1_TTL)

function get(key):
    if l1.allowed(key):
        v = l1.get(key)
        if v != null: return v
        v = coalesced_fetch(key)        # implementation 1
        l1.put(key, v, ttl=L1_TTL * uniform(0.8, 1.2))
        return v
    return coalesced_fetch(key)

Numbers: shard load for one key becomes hosts ÷ L1_TTL. With 500 hosts and a 2s TTL, a key at 1M reads/sec drops to ~250 reads/sec on its shard. Promotion latency is detection window + push: about 10–15s. Never allow-list a namespace marked zero-staleness; the client library should refuse. Detail: Distributed Cache, Redis.

3. Write Salting — Sub-Keys with a Discoverable Bucket Count#

# metadata (cached 60s): buckets["post:991:likes"] = 16; default 1
function incr(key):
    n = buckets.get(key, 1)
    store.increment(f"{key}#{random(0, n-1)}", 1)

function read(key):
    n = buckets.get(key, 1)
    return sum(store.multi_get([f"{key}#{i}" for i in 0..n-1]))   # cache the sum 1-5s

# Raising n is safe (new buckets start at 0). Lowering n requires
# draining the high buckets into the low ones first: a migration.

Numbers: on DynamoDB, pick n ≈ peak_writes ÷ 500 to keep each partition at 50% of its 1,000 WCU/s ceiling; a 16-way salt lifts a key to ~16K writes/s with a 16-way read fan-in. On Redis, salted counters only help if the sub-keys land on different shards, so use distinct hash tags. Detail: Write-Heavy Systems, DynamoDB, Leaderboard.

4. Pre-Aggregation — Collapse Writes Before They Reach the Key#

# Per host (or per stream subtask): aggregate for 1s, flush one write
local = Counter()
on like(post_id): local[post_id] += 1
every 1s:
    for (post_id, delta) in local.drain():
        store.increment(f"likes:{post_id}", delta)     # 1 write/s/host per hot key

# Stream variant (two-stage): stage 1 keys by (ad_id, hash(click_id) % 16),
# pre-aggregates per minute; stage 2 keys by ad_id and sums 16 partials.

Numbers: 38K INCR/sec from 200 hosts becomes 200 writes/sec. The cost is up to one flush interval of lost counts on a host crash, and a live count that is ~1s behind. For money, use the stream variant with exactly-once checkpoints rather than in-memory batching. Detail: Ad Click Aggregation, Stream Processing, Flink.

5. Split the Entity — Summary vs Detail, Global vs Regional#

Some keys are hot because they bundle too much. A 3 KB event summary that every viewer fetches includes a 2.8 KB detail block only 2% of viewers expand: split it and the hot key is 200 bytes, which moves the NIC ceiling 15× before any replication. A global leaderboard can become per-region boards plus a merged top-100 recomputed every 5 seconds. A 100K-member chat group moves from per-member inboxes (100K writes per message) to a shared group log with per-member cursors (Chat App). These are the cheapest fixes at runtime and the most expensive to ship, because each one changes the data model and usually the product contract. Detail: Schema Design.

6. Dedicated Cell — When the Entity Is Its Own Business#

When one tenant is 30–40% of traffic, or one account's events must never share fate with anyone else's, move it out: a directory entry routes that entity to its own shard group, its own cache cluster, sometimes its own deploy ring (Database Sharding, Partitioning). This is the only fix that eliminates neighbor impact rather than reducing it. It costs a cell's worth of fixed infrastructure (often $10–50K/month), a routing layer that every request consults, and a cell that drifts from the main fleet's configuration unless deploys go to it first.

Diagram: 6. Dedicated Cell — When the Entity Is Its Own Business

Architecture Diagram#

Diagram: Architecture Diagram

How to narrate it: the client library is where hot keys are handled, because it is the only component that sees every request and can act in-process. The control plane decides; the clients act. The tier registry and event calendar are where product decisions enter the system, as data rather than code.


Celebrity Accounts Are a Product Decision#

The hardest hot keys are people. A 100M-follower account is simultaneously read-hot (its profile), write-amplification-hot (fan-out of every post), write-hot (likes on its posts) and notification-hot (pushes to its followers). Each subsystem can patch its own symptom, and the case studies linked above each do exactly that. The Staff move is to notice the patches are four answers to one question: what experience does a celebrity's audience get, and what experience does the celebrity get?

DecisionOptionsWho DecidesTypical Default
Feed delivery for their postsPush to all followers / pull at read time / hybridProduct, with infra cost inputPull above ~1M followers; push below (News Feed)
Freshness of counts on their postsExact live / 1–5s stale / rounded ("1.2M")ProductRounded display, 2–5s stale
Notification fan-outImmediate to all / staged waves / digestProduct + notifications teamStaged waves over 5–15 min (Notifications)
Tier entry and exitFollower count / traffic / manualProduct, with hysteresisEnter at 1M, exit below 800K, reviewed weekly
Feature parityAll features / some features delayedProduct + legal (for verified/official accounts)Documented gaps, e.g. "edits propagate within 30s"

Why product has to own it. Each row trades user experience for cost and reliability, and the trade is visible to the most visible users on the platform. An engineer who silently sets the pull threshold at 1M has made a product decision without a product owner. When the celebrity's manager asks why their post took 20 seconds to appear for some fans, "the code has a threshold" is not an answer anyone wants to give.

Why the tier needs hysteresis. The News Feed celebrity fan-out incident is an account crossing 1M followers during a viral moment while the classifier read a daily cached count: still classified "push," 7.2M inbox writes queued. Tier membership must be evaluated on fresh data, enter and exit at different thresholds, and include a rate-based trigger (any author whose fan-out exceeds N writes/min is reclassified automatically) so a surprise celebrity is caught without waiting for the follower count.

🎯 Staff Move: "I want a celebrity tier as an explicit attribute, owned by product, rather than thresholds scattered through the feed, counter and notification code. Product signs off on what that tier gets: pull-based delivery, counts that are a few seconds stale, staged notifications. Infra's job is to make the tier cheap and to catch accounts that cross into it during a viral moment."


Failure Scenarios#

1. The Shared Shard — A Viral Counter Takes Out Unrelated Leaderboards#

A celebrity post's like counter sits on Redis shard 14 alongside 1.2M other keys, including another team's leaderboard.

t=0      Celebrity posts. likes:{post_9} goes from 50/s to 38K INCR/s.
t=+30s   Shard 14 command thread at 95% CPU. Other shards at 15%.
t=+45s   p99 for every key on shard 14: 0.3ms -> 40ms.
t=+60s   Leaderboard reads (another team) time out at 50ms; 3% of rank calls 5xx.
t=+2min  Client retries double the load on shard 14.
t=+4min  On-call runs key sampling, identifies likes:{post_9}.
t=+6min  Pre-aggregation path enabled for the post. INCR rate -> 200/s. Recovered.

Detection: redis.shard_cpu skew > 3× median; hotkey.top1_share > 5% of a shard's ops; client.retry_rate by shard. Blast radius: every key on shard 14, owned by at least four teams. Mitigation: enable local batching for the key; cap retries per shard at 10% of requests (retry budget). Prevention: write batching on by default for counter namespaces; counters for viral-capable entities isolated to their own shard group; retry budgets in the client library (Backpressure). Owner: counters team owns the batching path; cache platform owns shard isolation and retry budgets.

🎯 Staff Insight: The incident ticket goes to the leaderboard team, because their alert fired. The hot key belongs to someone else. Hot-key incidents are cross-team by construction, which is why detection and first-response mitigation belong to the platform, not the key owner.

2. The Partition That Couldn't Split — Today's Date as a Key#

A DynamoDB table keyed by PK = DATE#2026-10-03 records events. Writes are fine at 400/s in testing; launch-day traffic is 3,500/s.

t=0      Launch. Writes to DATE#2026-10-03 ramp to 3,500/s.
t=+20s   Throttling on one partition. Table-level consumed capacity: 12% of provisioned.
t=+1min  Dashboard: "capacity fine." On-call raises table capacity 2x. No effect.
t=+8min  Adaptive capacity and split attempts help briefly; all writes share one key.
t=+15min Throttled writes retried by SDK; client latency p99 4s; event loss at queue overflow.
t=+40min Hotfix: PK = DATE#2026-10-03#<0..9>. Readers updated to fan out 10 ways.

Detection: ThrottledRequests > 0 while ConsumedWriteCapacityUnits is far below provisioned; contributor insights top partition key. Blast radius: all writes for the day; downstream analytics missing 15 minutes of events. Mitigation: write-shard the key (10 buckets ≈ 5,000 writes/s at 50% headroom); buffer writes in a queue rather than dropping them. Prevention: design review checklist rejects low-cardinality or monotonic partition keys; load test at 3× forecast on the single hottest key, not on the table. See DynamoDB and Scaling Writes. Owner: service team owns key design; data platform owns the review checklist.

3. The Promotion That Flapped — Automatic L1 Without Hysteresis#

A team ships automatic hot-key promotion with a single threshold of 20K reads/s.

t=0      Key product:4471 at 22K reads/s fleet-wide. Promoted to L1.
t=+10s   Shard sees ~250 reads/s for the key. Sampler (measuring shard-side load) reports 250/s.
t=+20s   Below threshold. Demoted.
t=+30s   Shard load back to 22K/s. Promoted again.
t=+5min  30 promote/demote cycles. p99 oscillates 2ms <-> 35ms every 20s.
t=+12min Engineer pins the key manually. Root cause found: sampler measured post-L1 load.

Detection: hotkey.promotions_per_key > 3 in 5 min; latency oscillation with a fixed period. Blast radius: the key's readers plus its shard neighbors, every 20 seconds. Mitigation: manual pin with a 24h expiry. Prevention: sample at the client before L1 (demand, not served load); separate promote/demote thresholds (20K / 5K) with a 5-minute dwell; alert on flapping. Owner: cache platform team.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Read-hot key saturates shard CPU/NICshard.cpu or shard.net_bytes skew > 3× median; top-K names keyAll keys on the shardAuto L1 promotion; key replicationCache platform (detection), service (namespace staleness)
Write-hot counter on one thread/partitionstore.throttled_writes with low table utilization; hotkey.write_rateKey's writers + shard neighborsLocal batching; saltingOwning service team
Invariant-hot row under a saledb.row_lock_wait_ms p99; OCC conflict rate > 5%The sale and the cluster hosting itAdmission control; bucketed inventoryCommerce team + platform
Celebrity fan-out stormfanout_queue_lag_p99 > 30s; fanout_writes_per_author top-KAll authors' deliveryReclassify to pull; tiered queuesFeed team; product owns tier
Ordering-hot Kafka/Flink keyOne partition's lag ≫ others; one subtask busy ratio ~100%That partition's consumersSub-key where order allows; bigger consumerStream platform + producer team
Promotion flappinghotkey.promotions_per_key > 3 / 5 minKey + neighbors, periodicPin with expiry; hysteresisCache platform
Retry amplification on hot shardclient.retry_rate per shard > 10%Shard, then callers' thread poolsRetry budgets; circuit breaker per shardClient library owners (Stopping Cascading Failures)

Beyond Staff: The Principal View#

Why L7 Sees This Problem Differently#

A Staff engineer fixes the hot key in front of them. A Principal engineer notices that the same celebrity broke the feed, the like counter, the notification fan-out and the search index on the same evening, that four teams each built a different detector, and that the post-mortems all ended with "add the key to an allow-list." At org scale, hot keys stop being a data-structure problem and become a skew policy: which entities the company treats as special, who decides, what they cost, and which mitigations every shared-store client applies without anyone asking. The question is no longer "how do I fix this key" but "why does every team rediscover this during an incident, and what would make the cheap fixes automatic and the expensive ones deliberate?"

🧭 Principal Move: "We've had six hot-key incidents this year across four teams, and every one was fixed by hand-copying a key into an allow-list. I'd make coalescing, top-K detection and L1 promotion defaults in the shared client, create a product-owned tier registry for celebrities and whales, and require a design review only for the expensive fixes: salting, splitting, dedicated cells."

The Org-Level Fault Line#

Platform-owned automatic mitigation vs per-team hot-key handling.

OptionWhat WorksWhat BreaksWho Pays
Each team handles its own hot keysFixes fit the domain; no platform dependencyDetection reinvented 10 times; neighbors on shared shards unprotected; every team learns during its own incidentOn-call of whichever team shares the shard
Central platform does everything, including salting and splittingOne expert team, consistent behaviorPlatform can't see domain semantics; invariant-hot keys get "fixed" unsafelyProduct correctness; platform team becomes a bottleneck
Platform owns ranks 0–3 automatically; domain teams own ranks 4–7 with reviewCheap fixes everywhere by default; semantic fixes stay with people who understand the invariantNeeds a staleness declaration per namespace so the platform knows where L1 is allowedPlatform headcount (2–4 engineers); product time to declare tiers

The Principal default is the third row. The dividing line is semantics: anything that only adds copies or collapses duplicate reads is safe to automate; anything that changes what a key means (salting, splitting, moving) needs the owning team.

Cost Model#

Assumptions: managed Redis node ~$500–1,000/month; a dedicated cell (cache + store + routing) ~$15–40K/month; fully loaded engineer ~$25K/month; incident cost estimated from revenue at risk per minute of degraded service.

ScaleHottest Key PeakWhat You RunInfra $/monthPeopleRough Monthly Total
Startup (5K QPS total)~2K/sCoalescing in the client; shard CPU skew alert~$0 extra0.1 FTE~$3K (attention, not infra)
Growth (300K QPS)~150K/s (one celebrity, one sale)Shared client with top-K + auto L1; salted counters; admission for sales~$5–10K (hot-key service, headroom on hot shard groups)1 FTE across cache platform~$35K
Large (5M QPS, 200 teams)~2M/s (global live event)Org-wide client defaults; tier registry; event calendar; 2–3 dedicated cells for whales~$80–150K (cells + 30% headroom on shard groups)3–4 person platform team; game days~$250K

The Principal observation: at large scale the biggest line is not the hot-key service, it's the dedicated cells and headroom. Each cell should be justified by a named entity's revenue or contract, reviewed yearly, and folded back into the main fleet when the entity cools off. Cells that outlive their reason are the most common silent cost in this pattern.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoor TypeReversibility Cost
Enabling coalescing in the clientTwo-wayConfig flag
L1 TTL for a namespaceTwo-wayConfig change; product re-signs if raised
Salt bucket count for a keyOne-way-ishRaising is free; lowering needs a drain migration
Splitting an entity (per-region boards, summary vs detail)One-wayAPI and product semantics change; clients build on it
Celebrity tier semantics in a public API ("counts may be delayed")One-wayTightening later is a breaking change for consumers
Moving a whale to a dedicated cellTwo-way, expensiveWeeks of migration each direction; routing layer stays
Low-cardinality partition key in a new tableOne-wayFull table rewrite to change the key

The Standard I'd Write#

RFC: Hot-Key Handling Standard (v1)

Scope: Every service that reads or writes a shared cache, KV store, or stream through the platform client.

MUST:

  1. Use the platform client, which enables request coalescing and top-K key sampling (1% sample, 10s window) by default.
  2. Declare a staleness class per key namespace: S0 (no copies allowed), S2 (≤2s), S5 (≤5s). Automatic L1 promotion applies only to S2/S5 namespaces.
  3. Every design review names the hottest key, its forecast peak, and the ceiling of the partition it lands on. Peak > 50% of ceiling requires a documented fix from the ladder.
  4. Celebrity, whale and event entities are identified through the tier registry, not hard-coded thresholds.
  5. Salting, entity splits and dedicated cells require review by the data platform group and a named owner.

SHOULD: Register scheduled high-traffic events in the event calendar at least 24 hours ahead; run one hot-event game day per quarter on a top-10 service.

Exceptions: Filed with the platform team, time-boxed to 2 quarters, with an owner and expiry.

Success metrics: zero Sev-1s where the paged team didn't own the hot key; median time from heat to mitigation < 30s; number of manual allow-list entries trending to zero; every dedicated cell re-justified yearly.

What I'd Tell the VP#

When something goes viral, a celebrity posts or a big sale opens, all of that traffic lands on one small piece of our infrastructure, and it takes down other teams' features with it. This year that happened six times, and each time an engineer fixed it by hand. I'm proposing we build the standard defenses into the shared software every team already uses, so the system reacts in seconds on its own, and that product owns a short list of "special" accounts and events with agreed rules for how they're treated. It costs about two engineers for two quarters plus some reserved capacity. It pays back in fewer outages during exactly the moments we most want to be up. The one thing I need from product leadership is ownership of the celebrity and event list.

Principal Interview Signals#

SignalWhat It Sounds Like
Skew as a policy"Which entities does the company treat as special, and who decides? That's the question behind every hot-key incident we've had."
Automate by semantics"Anything that only adds copies is safe to automate. Anything that changes what a key means needs the owning team."
Product owns the tier"Celebrity behavior is a product contract, not a threshold in four codebases."
Pricing isolation"That whale's cell costs $30K a month. I want it re-justified yearly and folded back when they cool off."
Neighbor accountability"The team that got paged didn't own the key. I'd fix that with platform detection and per-shard retry budgets."

Staff answers that L7 interviewers find insufficient:

  • "I'd add an L1 cache for that key." Correct for one key; silent on why the next team will hit the same thing next month.
  • "Pull for accounts over 1M followers." A good threshold with no owner, no hysteresis, and no view of the counter and notification systems that break on the same night.
  • "Move the whale to its own shard." Isolation without a cost line or an exit plan.

How Real Companies Built It#

Facebook: Leases and Replication Within Pools (NSDI 2013)#

Facebook's memcache paper describes two hot-key mechanisms directly. Leases throttle refills: each memcached server hands out a lease token for a key at most once every 10 seconds by default, and other clients are told to wait briefly and retry. On a set of keys especially prone to thundering herds, the peak database query rate fell from 17K/s without leases to 1.3K/s with them. For keys whose request rate exceeds one server, they replicate a category of keys within a pool rather than splitting the key space further, and they note that a single key can account for 20% of a server's requests, which is why failed servers' load goes to a small dedicated Gutter pool (about 1% of servers) instead of being rehashed onto healthy ones (NSDI '13 paper).

Staff insight: This paper covers ranks 1 and 3 of the ladder, with numbers. The Gutter detail is the neighbor-protection argument: rehashing a hot key's load onto a healthy server just moves the incident.

Discord: Hot Partitions, Then Coalescing in a Data Service#

Discord's 2023 write-up on storing trillions of messages describes hot partitions in their Cassandra cluster: one channel-and-bucket pair receiving heavy traffic raised latency on its nodes, and every query to those nodes suffered. Part of the fix was a layer of Rust data services between the API and the database that coalesce requests (the first request for a piece of data spins up a task; concurrent requests subscribe to it) and use consistent-hash routing so all requests for one channel reach the same service instance. Alongside the move from 177 Cassandra nodes to 72 ScyllaDB nodes, p99 for fetching messages went from 40–125 ms to about 15 ms (Discord blog).

Staff insight: Routing all requests for one hot entity to one coalescing instance is a deliberate trade: it concentrates the key so it can be deduplicated before it reaches storage. A good reference for "coalescing is a design layer, not just a library call." The Chat App case covers the data model behind it.

Amazon DynamoDB: Adaptive Capacity and the Limits of Splitting#

DynamoDB documents a per-partition ceiling of 3,000 read units and 1,000 write units per second, and recommends write sharding for keys that exceed it. The 2022 USENIX ATC paper names hot partitions as one of the two most common customer problems, and reports that adaptive capacity eliminated over 99.99% of throttling due to skewed access. It also describes split for consumption: a partition that crosses a throughput threshold is split at a point chosen from observed key distribution. The paper is explicit that a partition receiving high traffic to a single item will not benefit, and DynamoDB detects that pattern and avoids splitting (partition key design, ATC '22 paper).

Staff insight: The managed service solves skew across keys and hands single-key heat back to you. Quote this when an interviewer says "DynamoDB will auto-scale." Details in DynamoDB.

Count-Min Sketch: Top-K Detection in Fixed Memory#

Cormode and Muthukrishnan's count-min sketch summarizes a stream in sublinear space and answers point queries with bounded overestimate; the paper applies it directly to finding frequent items (heavy hitters). This is the mechanism behind client-side hot-key detection: a few kilobytes per host give frequency estimates good enough to identify any key above a fraction of a percent of traffic. Redis ships a server-side counterpart, redis-cli --hotkeys, which only works when the eviction policy is LFU (count-min paper, redis-cli docs).

Staff insight: Naming the data structure, its memory and its error bound turns "we'd monitor it" into a design. The LFU requirement is a useful gotcha: many caches run LRU, and the built-in tool silently doesn't apply.


Practice Drill#

Prompt: "We run a live-sports app. During last month's final, one match's score key went from 5K to 600K reads per second, the Redis shard holding it hit its NIC limit, and the 'like' button on the same match drove 50K increments per second into one counter. Score, likes, and an unrelated chat feature on the same shard all degraded for 9 minutes. Design the fix."

Staff Answer

Three keys, three kinds of hot. The score key is read-hot and changes every few seconds, so 1–2s of staleness is acceptable (product confirms; the broadcast itself is delayed more than that). First, split the entity: the score payload is ~4 KB because it bundles match stats; the hot part is ~150 bytes (score, clock, status). At 600K reads/s that's ~90 MB/s instead of ~2.4 GB/s, well under a 10 Gbps NIC. Second, the key goes into the event calendar and is pre-promoted to L1 (TTL 1s) on all 400 app hosts 15 minutes before kickoff, so the shard sees ~400 reads/s for it; coalescing ensures one refresh per host. Score updates are pushed over WebSockets for clients that hold a connection (Push vs Poll), which removes most polling anyway. The like counter is write-hot and commutative: batch per host for 1s and flush one INCRBY, so 50K/s becomes ~400 writes/s; the displayed count is 1–2s behind, rounded ("48.2K"). The chat feature was a neighbor: it moves to its own shard group, and the client gets per-shard retry budgets (≤10% of requests) so timeouts on one shard don't double its load. Detection: client top-K sampling with promote at 20K/s and demote at 5K/s over 5 minutes; alerts on shard.net_bytes > 70% of NIC and shard.cpu skew > 3×. Validation: replay last final's traffic at 2× against staging. Owners: the match team owns the entity split and calendar entries; cache platform owns promotion, retry budgets and shard isolation.

Why this is L6:

  • Classifies each key (read-hot, write-hot, neighbor) and applies a different rung of the ladder to each, with the cost stated.
  • Does the bytes math: the NIC, not the CPU, was the ceiling, so shrinking the value is the first fix.
  • Uses the calendar for a scheduled event instead of relying on detection, and names owners for each piece.

What L7 adds:

  • Notices that every live event (other sports, concerts, elections) has the same shape and proposes the event calendar and tier registry as an org-wide standard, not a sports-team fix.
  • Prices it: pre-promotion and batching cost almost nothing; a dedicated shard group for live events is ~$8K/month against 9 minutes of degraded service during the highest-revenue hour of the season.
  • Sets up a quarterly hot-event game day replaying the biggest event of the previous quarter at 2×.

Staff Interview Application#

How to Introduce This Pattern#

"Before I go further, the hottest key in this design is [the celebrity profile / the sale SKU / the live match], at roughly [N] per second at peak. One shard handles about [M], so sharding alone won't save it. It's [read-hot], so I'll use coalescing and a short-TTL L1, which the client library turns on automatically when it detects the key. If that's not enough, I'll replicate the key. The one thing I won't do is split an inventory row without admission control in front."

Lead with the number, then the classification, then the cheapest fix that works, then what you'd escalate to and who signs off.

When NOT to Use This Pattern#

Don't add hot-key machinery (L1 promotion, salting, replication, cells) when:

  • The hottest key is comfortably under the ceiling. If the forecast peak of the single hottest key is under ~30–50% of one partition's ceiling, coalescing and a skew alert are enough. Salting a key that does 200 writes/s adds a 16-way read fan-in for nothing.
  • The data can't be stale and the load fits one node. Balances, permissions after revoke, inventory at the final decrement. Copies add correctness risk. Scale that partition vertically or add replicas of that shard with versioned reads.
  • Access is genuinely uniform. Batch exports, crawls and IoT telemetry with one write per device per minute have no hot key. Scale the store (Write-Heavy Systems).
  • The heat is really unbounded demand for a scarce thing. A 600K-person rush for 20K seats is not a hot key to absorb; it's demand to refuse. Use admission control (Ticket Drops, Concurrency Control), then deal with whatever residual heat reaches the store.
  • Order matters across the whole key. Salting a Kafka key or a Flink key breaks the per-key ordering you chose it for. Fix the key's scope instead, or give that partition a bigger consumer.

Follow-Up Questions to Anticipate#

Interviewer AsksWhat They Are TestingHow to Respond
"What if the hot key changes every minute?"Detection latency"That's why detection is automatic and client-side: a 10s window plus a push is ~15s to promote. For known events I pre-mark keys so there's no detection lag at all."
"How do you invalidate a key replicated 8 ways?"Write-side cost of read fixes"Invalidation fans out ×8, so I version the key and bump one version entry that readers check, rather than deleting 8 copies and racing refills."
"Salting means reads fan out. Isn't that worse?"Pricing the fix"For a write-hot counter, yes, reads go from 1 to 16 lookups, so I cache the sum for 1–2s. The read side was never the bottleneck."
"What about a hot key in Kafka?"Ordering vs throughput"If order matters only per user within the chat, I key by user-in-chat and the partition spreads. If order matters for the whole chat, I keep one partition and size its consumer, because salting breaks the guarantee."
"Who decides an account is a celebrity?"Product ownership"Product, through a tier registry with hysteresis. Infra adds a rate-based trigger so an account that goes viral is reclassified before the daily follower count catches up."
"Can't we just give the hot shard a bigger machine?"Vertical vs structural"For a 2× overshoot, yes, and it's the cheapest fix. For a 10× overshoot a bigger box doesn't exist, and the neighbors are still at risk."

Scorecard#

DimensionSenior (L5)Staff (L6)Principal (L7)
FramingShards and assumes even loadNames the hottest key, its rate and the partition ceiling up frontTreats skew as an org policy: tiers, calendar, defaults
ClassificationOne fix for all heatRead-hot vs write-hot vs invariant-hot, fix per classDraws the automation line at semantics: copies automatic, meaning changes reviewed
DetectionLogs and dashboards after the factSampled top-K, skew alerts, hysteresisPlatform contract; time-to-mitigation as an org metric
CostNot discussedStaleness, fan-in and N× writes named per fix$/month for headroom and cells, with yearly re-justification
OwnershipStore ownerService owns key design; platform owns detectionProduct owns tier policy; platform owns automatic mitigation; review board for cells

Strong Hire Signals

SignalWhat It Sounds Like
Names the key early"The hottest key here is the live match score, about 600K reads per second."
Separates read from write heat"The profile is read-hot, the counter is write-hot; different fixes."
Does the bytes math"600K × 4 KB is about 19 Gbps, so the NIC is the ceiling before the CPU is."
Protects neighbors"The team that gets paged won't own the key, so detection and retry budgets belong to the platform."

Lean No-Hire Signals

SignalWhy It Misses the Bar
"Consistent hashing will balance it"Confuses key placement with per-key load
Caches a write-hot counterDoesn't reduce the write rate at all
Splits invariant state without discussing correctnessShips oversells or stranded inventory

Common False Positives: Knowing count-min sketch math ≠ having a mitigation plan. Quoting DynamoDB's per-partition limits ≠ designing a key that stays under them. "We'd add a CDN" ≠ handling hot keys that require identity or freshness.


Capacity Planning Quick Reference#

Sizing the Hottest Key#

hot_key_peak      = top_entity_baseline × event_multiplier       # e.g. 5K/s × 120 = 600K/s
overshoot         = hot_key_peak / partition_ceiling              # > 0.5 means you need a fix
nic_limit_qps     = nic_bytes_per_sec / value_bytes               # 10 Gbps / 3 KB ≈ 400K/s
l1_shard_qps      = hosts / l1_ttl_seconds                        # 500 hosts / 2s = 250/s
salt_buckets      = ceil(peak_writes / (0.5 × partition_write_ceiling))   # DynamoDB: peak / 500
batched_writes    = hosts / flush_interval_seconds                # 200 hosts / 1s = 200/s
replica_read_cap  = n_replicas × per_shard_qps                    # 8 × 150K = 1.2M/s

Key Numbers Worth Memorizing#

NumberContext
~100–200K ops/sSimple commands per Redis shard (single command thread)
3,000 / 1,000 per sDynamoDB read / write units per partition
~10 MB/sConservative Kafka per-partition planning throughput
~400K/sReads of a 3 KB value that saturate a 10 Gbps NIC
50%Partition-ceiling share at which a key needs a deliberate fix
1–5 sL1 TTL for promoted hot keys
20K / 5K per sExample promote / demote thresholds (hysteresis)
0.1–1%Client-side key sampling rate
32 KBCount-min sketch with w = 2,048, d = 4 (~0.13% error bound)
17K → 1.3K/sFacebook's published DB peak without vs with leases

Common Pitfalls Checklist#

  • The hottest key and its forecast peak are named in the design doc, with the partition ceiling next to it
  • Each hot key is classified read-hot, write-hot or invariant-hot before a fix is chosen
  • Coalescing is on by default; L1 promotion is automatic with hysteresis and sampled before L1
  • Zero-staleness namespaces are excluded from L1 by the client, not by convention
  • Salt bucket counts live in discoverable metadata; lowering them has a migration plan
  • Scheduled events pre-mark their keys; celebrity and whale tiers live in a registry product owns
  • Retry budgets per shard stop timeouts from doubling load on the hot node
  • Every dedicated cell has a named entity, a cost line and a yearly review
  1. Loading the index…