Partition Count Immutability
also called Fixed Partition Count, Key-to-Partition Remap Hazard
The partition count of a keyed log is effectively permanent, because raising it changes which partition a key hashes to and therefore breaks per-key ordering and every table rebuilt by replaying that log.
A consumer group's lag is growing and one partition is the cause. The obvious fix is more partitions, the command takes a second, and it is available in every management console. It is also one of the few changes to a streaming platform that cannot be undone and that breaks correctness rather than performance.
The mechanism is arithmetic. A producer places a keyed record with murmur2(key) % numPartitions. Raise the
count from 12 to 24 and the result changes for roughly half of all keys. Records for key account-7 written
before the change sit in partition 3; records written after sit in partition 19. Two consumers now process
one key's history concurrently, in no particular order, and nothing in the platform reports a problem.
Partition counts also cannot be reduced at all in Kafka, so the decision is one-directional as well as permanent.
Why it matters
Per-key ordering is the guarantee most streaming designs quietly depend on. A cancel after a create, a balance after a debit, a delete after an upsert: all of them are correct only because one partition is read by one consumer in offset order. The partition count is therefore part of the data contract, not part of the capacity configuration, even though it lives in the same place as the capacity settings.
The damage is also retroactive. Any table rebuilt by replaying the topic from the beginning will read the pre-change and post-change records through a mapping that no longer matches how they were written.
Implementation patterns
- Size from consumer parallelism and over-provision 2×. Partitions are cheap to hold and expensive to add, so the asymmetry justifies deliberate waste. A topic needing 195 consumers today should be created with around 512 partitions.
- Use a custom partitioner that maps keys through a stable function such as consistent hashing over a fixed virtual-node space, when growth is genuinely unpredictable. This costs a non-default partitioner in every producer, in every language the company uses.
- Add a new topic rather than partitions. Create
orders.v2with the new count, dual-publish, migrate consumers, retire the old topic. This is slower and it is reversible at every step. - Never raise the count of a compacted topic. Compaction is per partition, so after a change the latest value for a key can sit in the old partition while a newer tombstone sits in the new one.
- Alarm on partition-count changes in the audit log, with the same seriousness as a schema change.
Industry example
This is documented Kafka behaviour rather than one company's story: the default partitioner's modulo arithmetic has been in place since the project's early releases in 2011, the inability to decrease partition counts is stated in the project's own documentation, and the standard guidance from Confluent has long been to over-provision because expansion is unsafe for keyed topics in production.
Capacity numbers for the other side of the trade are also documented. ZooKeeper-backed clusters were kept under roughly 4,000 partitions per broker and 200,000 per cluster because controller failover had to load full cluster state; KRaft has been demonstrated well beyond that at the metadata layer, with 10,000 to 20,000 per broker still a sensible practical bound on broker I/O grounds. Over-provisioning by 2× sits comfortably inside those limits for almost every workload.
Failure scenarios
- Out-of-order processing for existing keys, starting at the moment of the change and lasting as long as both partitions hold unprocessed records for a key. Downstream symptoms are a cancelled order showing as open or a stale balance that corrects itself on the next event.
- A compacted table rebuild that produces resurrected deletes, because a tombstone and an older value for one key live in different partitions with independent compaction.
- A stateful job whose keyed state no longer matches its input assignment, so a task holds state for keys it no longer receives and starts fresh state for keys it now owns. State size roughly doubles and aggregates restart.
- No improvement in the original problem. If lag was caused by one hot key, that key still hashes to one partition. The expansion added partitions and solved nothing, which is the most common outcome.
Trade-offs
| Choose | Gains | Pays |
|---|---|---|
| Over-provision partitions at creation | Expansion is never needed | More broker file handles, longer leader elections on broker loss |
| Custom stable partitioner | Safe expansion later | A non-default partitioner in every producer and language |
| New topic and migrate | Safe and reversible | Dual publishing, consumer migration, two topics' storage for a period |
When not to use it
The constraint only binds on keyed topics. If a topic carries independent records with no per-key ordering requirement, partition expansion is safe and routine, and treating it as a one-way door is false caution. Logs, metrics and most telemetry fall here.
It also stops binding when the downstream store is an upsert keyed by a business key with a monotonic version or timestamp, because out-of-order delivery is then resolved by the sink and ordering was never load-bearing. Teams that build this way buy themselves the freedom to repartition at will, and it is worth doing deliberately.
Interview question
Q: A consumer group is lagging and an engineer has already run the command to raise the topic from 12 to 48 partitions. The topic is keyed by account id and feeds a balance projection. What do you do in the next hour, and what do you tell the team afterwards?
What a strong answer covers: first, stop producing and establish whether any key has records in both the old and new partitions, because that is the window in which the projection can be wrong. Then: rebuild the projection from a point before the change using a single-threaded replay, or from a snapshot plus ordered re-application per key. Say plainly that the balance figures between the change and the repair cannot be trusted. Afterwards: the lesson is not "be careful with that command" but that the partition count is part of the contract, so it belongs behind the same review as a schema change, and that the original lag was almost certainly a hot key, which more partitions never fixes.
Quick check
Quiz: Why does raising a keyed topic's partition count break ordering for existing keys? Because key
placement is hash(key) % numPartitions, so a key's old and new records live in different partitions and are
consumed concurrently.
Flashcard: A lagging consumer group has one hot partition. You double the partition count. What improves? Nothing: the hot key still hashes to a single partition, and you have permanently broken per-key ordering for about half the key space.