skip to content

A team keys Kafka records by customer_id and a few large customers cause severe partition skew. Walk through the trade-offs and options for fixing it.

level: principalimportance: should knowfreq 34%

answer

  1. one hot key = one partition = one consumer ceiling
  2. ordering vs parallelism is the core trade
  3. composite/sub-key spreads a whale
  4. custom partitioner isolates whales
  5. adding partitions doesn't fix a single hot key

basics

~20 s

Keying by customer_id keeps each customer ordered but routes all of a hot customer's traffic to one partition, overloading it and its consumer. Options: composite keys to spread, a custom partitioner, more partitions, or accepting weaker per-key ordering — each trades ordering against balance.

solid answer

~50 s

The root cause: `murmur2(customer_id) % N` is deterministic, so a whale customer always hits one partition — and Kafka caps that partition's throughput at one consumer in the group, so you can't parallelize it. The tension is ordering vs. balance: strict per-customer ordering *requires* one partition per customer's stream. Options: (1) **Composite key** like `customer_id + bucket` (sub-key) to split a hot customer across K partitions — you keep ordering only within a sub-stream, so you must not need strict global per-customer order. (2) **Custom Partitioner** isolating whale keys onto dedicated partitions. (3) **Add partitions** — helps only if skew is spread across many keys, not a single hot one, and re-maps existing keys. (4) **partitioner.ignore.keys / null keys** if you don't actually need key ordering. (5) Rethink whether ordering is required at all, or move ordering downstream (e.g. sequence numbers + reordering). The decision is fundamentally: how much ordering can the business relax to gain parallelism.

go deeper

for a junior

Recognize that keying a whale to one partition overloads it and that one partition is read by only one consumer.

for a middle

Explain the ordering-vs-balance trade and that adding partitions doesn't help a single hot key.

for a senior

Lay out concrete options (composite key, custom partitioner, repartition) with their ordering costs.

for a principal

Drive from the actual ordering requirement, weigh migration/repartition hazards, and design monitoring + a migration plan, not just a quick fix.

**Why this happens.** Keying routes via `murmur2(customer_id) % numPartitions`, which is deterministic — every record for a given customer lands on exactly one partition. That's the whole point: it guarantees per-customer ordering and lets you process a customer's events in sequence. But it means a **single hot key cannot be split**: all of whale-customer's traffic funnels into one partition. **Why one hot partition is a real ceiling.** Within a consumer group, **each partition is consumed by at most one consumer**. So a hot partition is bottlenecked by a single consumer thread no matter how many you add. Skew also bloats that partition's log, lag, and replication load. You cannot fix it by scaling consumers alone. **The core trade-off: ordering vs. parallelism.** Strict per-customer ordering *demands* a single partition per customer-stream. Any spreading sacrifices some ordering. So every option below is really a question of how much ordering you can relax. **Options.** 1. **Composite / sub-keyed key.** Use `customer_id:bucket` where bucket = `hash(record) % K`. The whale now spreads across up to K partitions. You retain ordering *within a (customer, bucket)* sub-stream but lose strict global ordering for that customer. Works when ordering only needs to hold per smaller unit (e.g. per session, per device) rather than per whole customer. 2. **Custom Partitioner with whale isolation.** Detect the few hot keys and pin each to its own dedicated partition (or a salted set), while normal keys hash as usual. Keeps ordering, contains blast radius, but adds operational logic and must adapt as whales change. 3. **Increase partition count.** Only helps when skew is *distributional* (many medium keys colliding), because more partitions reduce collisions. It does **not** help a single dominant key — that key is still one partition. And `% N` re-maps existing keys, perturbing ordering history. 4. **Drop key-based routing** (`partitioner.ignore.keys=true` or null keys). Maximizes balance via uniform sticky partitioning but abandons per-key ordering entirely — only viable if ordering isn't actually required. 5. **Move ordering out of partitioning.** Attach monotonic sequence numbers and reorder downstream, freeing you to balance freely. Adds consumer-side complexity and state. **Decision framing (principal level).** Start from the *ordering requirement*, not the infrastructure: Is global per-customer order truly needed, or only per-session/per-aggregate? If the latter, sub-keying is clean. If a handful of whales dominate, isolation buys the most with least disruption. Beware repartitioning: changing N reshuffles all keys and can momentarily violate ordering during the transition — plan a migration (drain, dual-write, or accept a reset point). Finally, monitor partition-level lag/throughput to detect skew early rather than after it pages you.

  • Why doesn't simply adding more partitions fix a single dominant hot key?
    A single key always hashes to exactly one partition regardless of total partition count, so the whale still lands on one partition. Adding partitions only helps when skew comes from many keys colliding; it also re-maps existing keys via the changed modulo.
  • If you switch to a composite key customer_id:bucket, what ordering guarantee do you lose and keep?
    You keep ordering within each (customer, bucket) sub-stream, but lose strict global ordering across all of that customer's records, since they now span multiple partitions consumed independently. Acceptable only if the business needs ordering at the sub-stream granularity, not whole-customer.
  • What ordering hazard does increasing the partition count introduce during the change itself?
    Because routing is hash % N, changing N re-maps keys to different partitions. In-flight and new records for a key may land on a different partition than its history, so per-key ordering can be violated across the cutover unless you drain or plan a reset point.

saying these in an interview costs you the question

  • Suggesting 'just add consumers' (one partition is still one consumer)
  • Claiming adding partitions fixes a single hot key
  • Ignoring that any spreading scheme weakens per-key ordering
  • Forgetting that changing partition count re-maps existing keys and can break ordering during migration

context