A partitioned topic is keyed by user ID to preserve per-user event ordering, but one 'celebrity' user generates 100x the event volume of a typical user. What problem does this cause, and what are reasonable ways to mitigate it without abandoning per-user ordering entirely?
answer
- hot key -> hot partition, always
- aggregate lag can hide per-partition lag
- sub-keying spreads load, breaks strict order
- isolate known hot keys onto dedicated resources
- consumer must reassemble order if sub-keyed
basics
~20 sThat one heavy user's events all pile onto a single partition no matter how many partitions the topic has, so that one partition gets overloaded and lags behind while the rest stay fine. Fixes usually involve splitting that user's events across a few partitions and reordering them again on the consumer side, or treating that user as a special case.
solid answer
~50 sBecause a partition key always maps to exactly one partition, a single high-volume key concentrates all of its traffic onto one partition regardless of the total partition count - this is 'key skew' or a 'hot partition,' and it shows up as one partition's consumer lag climbing while sibling partitions stay near zero, capping the whole group's effective throughput at what that one partition/consumer can sustain. Mitigations trade some ordering guarantee for load spreading: sub-keying (append a bounded random or round-robin suffix to the hot key, e.g., userId#0..userId#7, spreading it across a handful of partitions) sacrifices strict ordering for that user unless the consumer reassembles order using event sequence numbers or timestamps after merging; alternatively, detect hot keys and route them to a dedicated, isolated topic/partition sized for their load, or accept eventual/best-effort ordering for that entity specifically while keeping strict ordering for the long tail of normal-volume keys.
go deeper
Should recognize that one very active entity can overload a single partition even with many partitions total.
Should know per-partition lag monitoring is needed to detect this, since topic-level aggregate lag can hide it.
Should propose concrete mitigations (sub-keying, hot-key isolation) and articulate what ordering guarantee each one gives up.
Should design a monitoring-plus-mitigation strategy that automatically identifies emerging hot keys and applies isolation/sub-keying policy without requiring the entire system's partitioning scheme to be redesigned around a rare outlier.
## Why a hot key concentrates on one partition **Key skew** is the direct consequence of the same mechanism that provides per-entity ordering: since a partition key deterministically maps to one partition via `hash(key) % numPartitions`, every event for a given key is, by design, funneled to the same place. That's exactly what you want for ordering, but it means the partitioning scheme has no way to spread load within a single key's traffic — partition count only helps distribute load across different keys. If one key (a celebrity user's account, a mega-tenant in a multi-tenant SaaS product, a single stock symbol during a trading spike) generates disproportionate volume, all of that volume lands on one partition no matter whether the topic has 6 partitions or 600. ## The symptom, and why it hides The observable symptom is a **hot partition**: its consumer-side lag (the gap between the latest produced offset and the last committed offset) climbs steadily while every other partition in the same topic stays near zero, because the one consumer instance assigned that partition is saturated while its siblings are comfortably idle. This is dangerous precisely because aggregate lag metrics across the whole topic can look acceptable — if 23 of 24 partitions are healthy and only 1 is falling behind, a topic-level average can mask a user-facing problem (that one entity's downstream state is stale or delayed) entirely. Teams that only alert on aggregate lag, not per-partition lag, routinely miss this until a customer complains. ## Three mitigations, and what each gives up The trade-off in every mitigation is the same: you're relaxing the guarantee that all of one entity's events are strictly ordered on a single partition, in exchange for spreading that entity's load across more resources. | Mitigation | What it means for the hot entity | |---|---| | key salting/sub-keying | restores parallelism for the hot entity but breaks strict cross-sub-partition order | | isolating known hot keys onto dedicated infrastructure | keeps strict single-partition ordering for the hot entity | | increasing total partition count, when the skew is moderate rather than extreme | the hot key's ceiling is now whatever one partition/one consumer instance can sustain, as a known, monitored limit | 1. **A common technique is key salting/sub-keying**: append a bounded, deterministic-or-round-robin suffix to the hot key (e.g., userId becomes userId#0 through userId#7), producing 8 sub-streams for that one entity that spread across up to 8 partitions instead of 1. This restores parallelism for the hot entity but breaks strict cross-sub-partition order — the consumer side must now either tolerate approximate ordering (fine for, say, an analytics/metrics pipeline) or actively reassemble true order using an application-level sequence number or timestamp embedded in each event, essentially reimplementing ordering in the consumer instead of relying on the log for it. 2. **A second approach is isolating known hot keys onto dedicated infrastructure**: detect (via monitoring or a pre-known list, like verified/celebrity accounts in a social platform) which keys are disproportionately large and route them to a separate topic sized and partitioned specifically for that traffic pattern, or even give each known hot key its own dedicated partition via a custom partitioner. This keeps strict single-partition ordering for the hot entity while not making every other, normal-volume key suffer from a partitioner tuned around the outlier. 3. **A third, cheaper mitigation when the skew is moderate rather than extreme** is simply increasing total partition count and relying on a well-distributed key space for the rest of the traffic, while accepting that the hot key's ceiling is now "whatever one partition/one consumer instance can sustain" as a known, monitored limit — appropriate when the hot entity's absolute volume, while large relative to peers, is still within what a single fast consumer can handle (e.g., a well-optimized consumer doing simple filtering/routing may sustain far more throughput per partition than one doing heavy per-event computation). ## The failure mode to watch for The failure mode to watch for in all of these: sub-keying without a corresponding consumer-side reordering strategy silently reorders a hot entity's events for anyone downstream who assumed strict order, which can be worse than the lag problem it solved — e.g., an "account suspended" event processed before a "transaction" event from the same sub-key split, because they landed on different sub-partitions and were consumed at different rates, causing a transaction to post to a supposedly-suspended account. ## Where it shows up A concrete real-world scenario: ride-sharing and social-media platforms both commonly hit this with celebrity/power-user accounts or extremely active drivers/regions — a location-update or activity-event stream keyed by user ID can have one influencer account generating orders of magnitude more events than a typical user during a viral moment; production systems typically respond by combining per-partition lag alerting (to catch it early) with a sub-keying or hot-key-isolation strategy specifically reserved for identified outliers, rather than redesigning the partitioning scheme for the entire, otherwise well-behaved key space.
- How would you detect a hot-key problem in production before it causes a customer-visible incident?Monitor consumer lag per partition, not just aggregate lag across the topic, and alert on any single partition's lag diverging significantly from its siblings; pairing that with per-key volume metrics (if feasible) helps confirm the cause is a specific hot key rather than, say, a slow downstream dependency affecting one consumer instance.
- If you sub-key a hot user's events across 8 partitions, how can a downstream consumer still reconstruct that user's true event order when needed?Embed a monotonically increasing sequence number or a high-resolution producer timestamp in each event at produce time, then have the consumer buffer and merge events for that user's sub-keys (e.g., a short windowed buffer keyed by userId across all 8 sub-partitions) and emit them in sequence-number order, trading some added latency and consumer complexity for restored ordering.
It's like a bank that assigns each customer to one specific teller for continuity of service: a regular customer's transactions all flow through smoothly, but a customer running a busy small business through their personal account can single-handedly keep that one teller backed up all day while every other teller stands idle.
saying these in an interview costs you the question
- Assumes adding more partitions automatically fixes a hot-key problem
- Doesn't recognize that a single key can never span multiple partitions under standard hashing
- Only monitors aggregate/topic-level lag, never per-partition
- Proposes sub-keying without addressing how the consumer will restore order for entities that need it
- Treats hot-key skew as a rare/theoretical concern rather than a common production pattern