04 / 05

A topic has 100 partitions but one receives most traffic. How would you diagnose and fix it?

Difficulty: 8/10
Partitioning and ordering

Confirm it is data skew, find whether it is a hot key or a bad key distribution, then spread load at the finest ordering granularity the business allows

I would first confirm the symptom with numbers, then classify the cause, because the fixes are different. To measure, I compare log end offsets and on-disk size per partition and look at per-partition consumer lag. A big partition on one broker can also be a broker-level problem (leader imbalance), so I check leader distribution too. Then I sample keys landing in the hot partition: one dominant key is a hot key; a few keys dominating is low key cardinality; many keys hashing unevenly is rare with murmur2 and usually points to a key format problem or a custom partitioner bug; null keys with an explicit partition set by code is another common culprit.

javascript

The fix depends on what ordering the business really requires. Option one: use a finer key, for example tenantId plus userId instead of tenantId, keeping per-entity ordering while spreading load. Option two: salt only the hot keys, so a hot key is spread over N partitions and downstream aggregation recombines results; this gives up ordering for that key across the sub-streams. Option three: a custom partitioner that spreads known hot keys and hashes everything else normally. Option four: if ordering does not matter, drop the key and let the sticky partitioner balance. Adding partitions does not fix a hot key, since the key still maps to one partition, and it breaks key mapping for everyone.

javascript
  1. 1

    Trade-off: every technique that spreads a hot key gives up strict ordering for that key. State that explicitly and ask whether the real requirement is ordering per user, per session, or per aggregate.

  2. 2

    Consumer-side mitigation: if one partition's consumer is the bottleneck, process records from that partition concurrently with a key-aware worker pool, preserving order per key inside the consumer. This is a workaround, not a cure for storage and broker skew.

  3. 3

    Common mistake: assuming skew means the hash function is bad. murmur2 is well distributed, so the cause is almost always the key choice, not the hash.

  4. 4

    Common mistake: adding partitions or consumers as the first reaction. Neither moves traffic off a hot key.

  5. 5

    Every producer must use the same custom partitioner, or ordering and locality assumptions silently break. Plan the rollout and monitor partition balance after the change.

  6. 6

    Version note: Cruise Control and tooling for leader and replica balancing exist for broker-level skew, but their availability depends on your distribution or managed offering.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.