Bucketed Partitioner
also called Virtual Partition Map, Bucket-to-Partition Map
Hashing keys into a large fixed set of buckets and routing buckets to partitions through an editable map - so partition count can grow without re-mapping keys or breaking per-key ordering.
A topic keyed by buyer id has 24 partitions, one consumer each, and that is now the throughput
ceiling. Raising it to 96 looks like a configuration change. It is not: the default partitioner is
hash(key) mod N, and changing N re-maps most keys. Only about one key in four is congruent under
both 24 and 96, so for a window of time a buyer's events exist in two partitions with no ordering
relationship between them, and every table rebuilt by replaying the topic reads an interleaving
that was never produced. Nothing errors; a subset of keys is simply wrong.
A bucketed partitioner removes the modulus from the hot path. The key hashes to one of a large, fixed number of buckets - 4096 is a common choice - and a separate map says which partition each bucket is served by. Keys never change bucket, because the bucket count never changes. Growing the partition count becomes an edit to the map, applied one bucket at a time.
This is the same move a sharded database makes with a shard map, and the same reason consistent hashing uses virtual nodes: route through a table you can edit rather than through an arithmetic identity you cannot.
Why it matters
Partition count is otherwise close to permanent for any keyed topic, which means it has to be chosen years in advance - over-provisioned, at a real cost in broker metadata, open file handles, replication fan-out and rebalance time, or under-provisioned, with a ceiling nobody can raise.
The map also makes the dangerous operation incremental and reversible. A migration can move one bucket - about 1/4096 of keys - pause producers for those keys for a fraction of a second, confirm consumers have drained the old location, flip one row, and release. The blast radius of a mistake is a small share of traffic, and the rollback is flipping the row back.
Implementation patterns
- Fix the bucket count on day one and never change it. It must exceed any partition count you might reach, by a wide margin. 4096 against a plausible ceiling of a few hundred partitions is cheap.
- Store the map where producers and consumers both read it, versioned, with a change log. It is configuration data, so treat it as data: reviewed, auditable, revertable.
- Migrate per bucket. Mark draining, buffer the bucket's keys at the producer, wait for end-of-partition on the old location, flip, release.
- Never dual-produce the same key to two topics to smooth a cutover. That is precisely the ordering loss the pattern exists to prevent.
- Add a per-key sequence number to the event and assert consumer-side that it never goes backwards. Divergence then shows up within a minute rather than in a month-end reconciliation.
- Keep buckets balanced by count, not by key. Monitor bytes per bucket, because one heavy bucket reproduces the skew problem inside the new scheme.
Industry example
The need is visible in production systems rather than in any vendor feature: sharded databases route by a key range or a shard map so that splitting a shard moves a range instead of rehashing the world, and distributed hash tables assign many virtual tokens per node for the same reason. Kafka's own documentation is clear that adding partitions to a keyed topic changes key placement and that the platform cannot repair the ordering consequence - which is why teams that expect to grow build the indirection themselves, in the producer, before the first message is written.
Failure scenarios
- Bucket count too small. 64 buckets against 96 wanted partitions leaves buckets that cannot be split, so the ceiling returns.
- Map skew. Buckets assigned evenly by count while traffic is uneven, so one partition is hot and the scheme hides it behind an extra layer.
- Stale map on one producer. A deploy misses an instance, which writes a migrated bucket to its old partition. Per-key ordering breaks for that bucket until someone notices the version mismatch - which is why the map version belongs in a metric.
- The forgotten replay path. The new topic's history starts at cutover, so until a full retention period has passed, a rebuild must read the old topic to its end and then the new one, per key.
- One dominant key. A single buyer generating a large share of events still lands on one partition and one consumer. The pattern fixes growth, not skew.
Trade-offs
It buys a partition count that is no longer permanent, and an incremental, reversible migration path. It pays with a lookup and a piece of distributed configuration on the producer's hot path, a deployment coupling (every producer must agree on the map), and an operational concept that every on-call engineer now has to understand. Choose it when the topic is keyed, long-lived and plausibly growing. Skip it when the topic is unkeyed, because then ordering does not matter and partitions can be added freely.
When not to use it
Do not retrofit it to solve a throughput problem that is really per-record CPU. A key-sharded worker pool inside each consumer - hash the key to one of eight in-process lanes, one thread each - multiplies concurrency by 8 with no topology change, preserves per-key order and reverses with a config flag. 24 partitions times 8 lanes is 192-way concurrency, which is usually enough and is the cheaper answer.
And if the per-partition distribution is already dominated by a few keys, fix the key design first. Repartitioning is an expensive way to discover you have skew.
Interview question
Q: A keyed topic needs to go from 24 to 96 partitions with per-buyer ordering as a hard requirement and six downstream jobs that rebuild state by replay. How do you sequence it, and where is the point of no return?
What a strong answer covers: why hash mod N makes this unsafe and that roughly three keys in
four move; trying in-consumer key-sharded parallelism first; a new topic with a fixed bucket count
and an editable bucket-to-partition map; per-bucket drain-then-flip rather than dual-producing;
per-key sequence numbers to detect divergence in minutes; fixing the replay path to read old-then-new
in order until retention has turned over; and naming the point of no return as the first key
produced to the new topic while its old-topic events are still unconsumed.
Quick check
Quiz: Why does raising a keyed topic from 24 to 96 partitions break ordering? Because
hash(key) mod N re-maps about three keys in four, so a key's new events land in a different
partition from its history and the two have no ordering relationship.
Flashcard: What does a bucketed partitioner make editable that the default makes permanent? - The bucket-to-partition assignment, and therefore the partition count: growth moves buckets through a map instead of rehashing every key.