skip to content

Why does increasing the partition count of a keyed topic break ordering and key-locality guarantees, and how do you avoid that disruption?

level: seniorimportance: must knowfreq 65%

answer

  1. hash(key) % numPartitions
  2. change N -> key remaps
  3. old data stays, new data moves
  4. order only within a partition
  5. republish to a new topic to re-lay-out

basics

~20 s

The default partitioner routes a key via hash(key) % numPartitions. Changing numPartitions changes the modulo result, so a key that always landed in partition 2 may now land elsewhere — splitting that key's records across partitions and breaking per-key ordering.

solid answer

~50 s

Kafka's default partitioner places a keyed record by computing `murmur2(key) % numPartitions`. Per-key ordering and 'all events for a key land in one partition' depend on numPartitions being stable. When you `--alter` to add partitions, the modulus changes, so the same key can hash to a different partition than before. New records for that key go to the new partition while old records remain in the old one — now the key's history is split across two partitions, and since ordering is only guaranteed within a partition, a consumer can process the key's events out of order. To avoid this: (1) over-provision partitions up front so you never need to grow; (2) use a custom partitioner or explicit partition assignment that is independent of partition count; or (3) accept the disruption only at a controlled cutover and repartition-by-republish into a fresh topic so all of a key's data is re-laid-out consistently.

code

java · 3 lines
java
int partition = Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
// numPartitions = 6  -> murmur2("cust-42") % 6 = 2
// numPartitions = 12 -> murmur2("cust-42") % 12 = 8   // key moved!

go deeper

for a junior

Know that partition for a key = hash(key) mod partitions, so changing the count moves keys.

for a middle

Explain that old data stays put while new data moves, splitting a key across partitions.

for a senior

Connect the split to the per-partition-only ordering guarantee and prescribe republish or headroom.

for a principal

Design key-routing strategies (custom partitioner, sizing) that make topics resilient to growth, and own the cutover plan for stateful/Streams consumers.

## Background: how keyed records are placed When a producer sends a record **with a key**, the default partitioner decides the target partition deterministically: ``` partition = Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions ``` (`murmur2` is the hash function; older clients used a different but analogous scheme.) Two properties follow: 1. **Key locality** — the same key always maps to the same partition, *as long as numPartitions is constant*. This is what lets you say 'all events for customer 42 are in one partition.' 2. **Per-key ordering** — because Kafka guarantees order only within a partition, key locality is what gives you ordered processing per key. ## What breaks on increase `partition = hash(key) % N`. The `% N` term depends on N. Increase N from 6 to 12 and for most keys `hash(key) % 12 != hash(key) % 6`. Concretely: - Before: key 'cust-42' → partition 2 (all its events sit in partition 2, ordered). - After increasing to 12: 'cust-42' now → partition 8. New events go to partition 8; the old events still live in partition 2. - A consumer reading partition 2 and partition 8 has no cross-partition ordering, so it can see a later event (partition 8) before an earlier one (partition 2) for the same key. This silently corrupts any logic that assumes single-partition, ordered key streams: stateful aggregations, dedup, sequence-number checks, CDC ordering, etc. ## Note: increasing does NOT move old data The increase only changes routing for *new* writes. Existing records are never re-shuffled into the new partitions, so the split is unavoidable with an in-place increase. ## Avoidance strategies 1. **Provision headroom up front.** Choose a partition count generous enough that you never have to grow. The classic trick: pick a count with many divisors / a power-of-two-friendly number so future consumer scaling stays within the existing partitions. 2. **Decouple key routing from partition count.** Use a custom `Partitioner` that maps keys via a stable scheme (e.g. consistent hashing or an explicit key→partition table) so adding partitions doesn't remap existing keys. Or have producers set the partition explicitly. 3. **Controlled repartition-by-republish.** When you must change the layout, create a new topic and republish: read the old topic and produce into the new one. Because the producer re-keys everything against the new partition count in one pass, every key is consistently placed in the new topic. Cut consumers over to the new topic. This is the only way to also *reduce* partitions or change the partitioning key. 4. **Quiesce keys during cutover.** If a brief disruption is acceptable, drain in-flight processing for affected keys before/after the change to bound out-of-order risk. ## Edge cases - Records **without a key** are not affected by key locality (they're spread by the sticky/round-robin partitioner), so order-sensitivity is the keyed-topic concern. - Kafka Streams repartition topics are especially sensitive — changing their partition count invalidates state-store/partition alignment, which is why Streams pins and validates partition counts.

  • After increasing partitions, a key's events are split between the old and new partition. Why can a consumer now see them out of order?
    Kafka guarantees ordering only within a single partition. Once a key's events span two partitions, the consumer reads the two partitions independently with no cross-partition ordering, so a newer event in the new partition can be processed before an older event still in the old partition.
  • How does repartition-by-republish produce a consistent layout when an in-place increase doesn't?
    Republishing re-keys every existing record against the new topic's partition count in a single producing pass, so all of a key's history lands in the same partition. An in-place increase only re-routes new writes, leaving old records stranded in their original partitions.

saying these in an interview costs you the question

  • Saying an increase preserves which partition a given key maps to.
  • Claiming Kafka rebalances existing keyed records into the new partitions automatically.
  • Asserting ordering is global across the topic rather than per-partition.
  • Thinking keyless records suffer the same per-key ordering break (they don't carry that guarantee).

context