skip to content

Your keyed topic shows severe partition skew (a few partitions hold most of the load). What are your options, and what is the fundamental tension you must navigate?

level: seniorimportance: should knowfreq 50%

answer

  1. skew = key distribution, not the hash
  2. hot key / whale concentrates on one partition
  3. salt key#0..N spreads but breaks total order
  4. custom partitioner isolates hot keys
  5. more partitions won't fix a single whale
  6. ordering only within a partition = the core tension

basics

~20 s

Skew comes from uneven key distribution (hot keys), not the hash. Options: salt or compound the key, route hot keys with a custom partitioner, or increase partitions. The tension: spreading a key for balance destroys the per-key ordering the key was meant to guarantee.

solid answer

~60 s

Partition skew on a keyed topic almost always traces to **key distribution**, not the murmur2 hash — a few hot keys (a dominant tenant, a null-ish default) concentrate traffic. Options: (1) **Re-key / salt:** append a bounded suffix (e.g. `key#0..N`) so one logical key spreads across N partitions — but this breaks global per-key ordering; you only keep ordering within each salted sub-key. (2) **Custom Partitioner** to route specific hot keys to dedicated partitions while hashing the rest, bounding skew without disturbing cold keys. (3) **Increase partitions** — helps only if the skew is from many medium keys, not one whale, and it reshuffles existing key mappings. (4) **Compound key** to raise cardinality. The fundamental tension: Kafka gives ordering *only within a partition*, so the very thing that causes skew (all of a key's records on one partition) is also what guarantees that key's ordering. You cannot both perfectly balance a hot key and keep its total order — you must decide whether ordering is required per logical key or only per finer-grained sub-key.

go deeper

for a junior

Understand that skew means some partitions get much more load, usually because of popular keys.

for a middle

Explain that the hash is uniform so skew is a key-distribution issue, and that adding partitions won't fix one hot key.

for a senior

Lay out salting, custom routing, and cardinality options and articulate the ordering-vs-balance tradeoff.

for a principal

Frame the decision around the true ordering unit, design a hot-key isolation/sharding strategy, and own the consumer-side and resize implications.

## Where skew comes from **Partition skew** means some partitions carry far more data/throughput than others, overloading specific brokers and the consumers that own those partitions. On a **keyed** topic, the murmur2 hash is statistically uniform, so skew is almost never the hash's fault — it comes from the **key distribution**: - **Hot keys (a whale):** one tenant/user generates most events; all of them hash to one partition. - **Low cardinality:** few distinct keys can't spread across many partitions. - **Degenerate keys:** e.g. many records with the same default/empty key. ## The fundamental tension Kafka guarantees ordering **only within a partition**. To get per-key ordering you must keep *all* of a key's records on *one* partition. But that is exactly what creates skew for a hot key. So: > Balancing a hot key across partitions necessarily sacrifices that key's total ordering. You cannot have both perfect balance and total per-key order for the same hot key. Every remedy is a point on this tradeoff. ## Options 1. **Salting / sub-keying:** transform the key to `originalKey#s` where `s` is in `0..N-1` (chosen round-robin or by a secondary field). The whale now spreads over N partitions. **Cost:** you lose global ordering for that key — you only retain ordering within each `key#s` bucket. Acceptable when ordering is needed per session/sub-entity, not per whole key. Consumers must understand the scheme to re-aggregate. 2. **Custom Partitioner for hot keys:** detect known hot keys and route them to dedicated/reserved partitions (or a wider sub-range), while everything else hashes normally. Bounds skew and isolates the whale without disturbing the cold majority. Still subject to the ordering tradeoff if a single whale spans multiple partitions. 3. **Increase partition count:** only helps when load is spread across **many** medium keys; a single whale still lands on one partition no matter how many you add. Also, adding partitions reshuffles existing `hash % n` mappings, breaking ordering across the resize and not migrating old data. 4. **Raise key cardinality (compound key):** combine fields (tenant+region+entity) so keys spread, if your ordering requirement is at that finer grain. 5. **Throttle / shard upstream:** sometimes the right answer is to split the whale's workload at the application level. ## Diagnosing Look at per-partition message rate / lag (e.g. via consumer lag tools, `kafka-consumer-groups.sh`, JMX `MessagesInPerSec` per partition). Confirm it's a key-distribution issue before touching the partitioner. ## Decision frame Ask: *what is the true unit that must stay ordered?* If it's the whole tenant, you're stuck with the whale on one partition (scale the consumer vertically, or shard upstream). If finer ordering (per session, per order id) suffices, salt or re-key to that grain and reclaim balance.

  • Will simply adding more partitions fix skew caused by one dominant key?
    No. A single hot key hashes to exactly one partition regardless of how many partitions exist, so its load stays concentrated. Adding partitions only helps when skew comes from many medium keys, and it also reshuffles existing key-to-partition mappings.
  • What do you give up when you salt a hot key to spread it across partitions?
    Global per-key ordering. After salting, ordering is only preserved within each salted sub-key bucket; records of the original key are now interleaved across multiple partitions with no total order, so consumers must tolerate or re-aggregate that.
  • How do you confirm skew is a key-distribution problem rather than a hashing bug?
    Inspect per-partition throughput/lag (kafka-consumer-groups.sh, per-partition JMX MessagesInPerSec). murmur2 is statistically uniform, so concentration on specific partitions almost always means a few keys dominate the input distribution.

saying these in an interview costs you the question

  • Blaming murmur2 / the hash function for skew instead of the key distribution.
  • Claiming more partitions always fixes skew (useless against a single whale).
  • Salting a key while still promising total per-key ordering.
  • Forgetting that adding partitions reshuffles existing key mappings and breaks ordering across the resize.
  • Ignoring that ordering-vs-balance is a genuine tradeoff, not something you can fully optimize both ways.

context