3.9 Partitions
That's a genuine, deliberate availability-for-consistency trade.
9.1 What keys are for
"Kafka messages are key-value pairs, and while it is possible to create a
ProducerRecordwith just a topic and a value, with the key set to null by default, most applications produce records with keys."
Keys serve two goals:
- Additional information stored with the message.
- Deciding which partition the message is written to. (Keys also play an important role in compacted topics — Ch. 6.)
"All messages with the same key will go to the same partition. This means that if a process is reading only a subset of the partitions in a topic, all the records for a single key will be read by the same process."
// with key
ProducerRecord<String, String> record =
new ProducerRecord<>("CustomerCountry", "Laboratory Equipment", "USA");
// null key — just leave it out
ProducerRecord<String, String> record =
new ProducerRecord<>("CustomerCountry", "USA");9.2 Default partitioner behavior
Null key → round-robin, and since 2.4, sticky round-robin
When the key is null and the default partitioner is used, the record goes to one of the available partitions at random; a round-robin algorithm balances messages among partitions.
Starting in the Apache Kafka 2.4 producer, the round-robin algorithm for null keys is sticky. "This means that it will fill a batch of messages sent to a single partition before switching to the next partition. This allows sending the same number of messages to Kafka in fewer requests, leading to lower latency and reduced CPU utilization on the broker."
| Pre-2.4 round-robin | 2.4+ sticky round-robin |
|---|---|
m1→p0 m2→p1 m3→p2 m4→p0 … | m1, m2, m3… → p0 (fill the batch), then m… → p1 (fill the batch) |
| Many small batches → many requests. | Few full batches → fewer requests → lower latency and less broker CPU. |
Non-null key → hash
Kafka will hash the key — "using its own hash algorithm, so hash values will not change when Java is upgraded" — and use the result to map to a partition.
A crucial detail with a real consequence:
"Since it is important that a key is always mapped to the same partition, we use ALL the partitions in the topic to calculate the mapping — not just the available partitions. This means that if a specific partition is unavailable when you write data to it, you might get an error." (Fairly rare — see Ch. 7 on replication and availability.)
| Key | Maps over | Consequence |
|---|---|---|
| Null key | available partitions | never blocked by an offline partition |
| Non-null key | all partitions | can error if the target partition is offline — necessary, or the key→partition mapping would drift every time a broker hiccuped |
That's a genuine, deliberate availability-for-consistency trade. You cannot have both stable key routing and immunity to offline partitions.
9.3 Other built-in partitioners: RoundRobinPartitioner and UniformStickyPartitioner
These give random and sticky random partition assignment even when messages have keys.
Why you'd want that:
"These are useful when keys are important for the consuming application (for example, there are ETL applications that use the key from Kafka records as the primary key when loading data from Kafka to a relational database), but the workload may be skewed, so a single key may have a disproportionately large workload. Using the
UniformStickyPartitionerwill result in an even distribution of workload across all partitions."
This decouples two things people conflate: the key as data vs the key as routing. You can keep the key for downstream semantics while abandoning key-based routing.
9.4 ⚠️ Adding partitions breaks key→partition mapping
*"When the default partitioner is used, the mapping of keys to partitions is consistent only as long as the number of partitions in a topic does not change. As long as the number of partitions is constant, you can be sure that records regarding user 045189 will always get written to partition 34. This allows all kinds of optimization when reading data from partitions.
However, the moment you add new partitions to the topic, this is no longer guaranteed — the old records will stay in partition 34 while new records may get written to a different partition."*
- · ordering broken
- · stateful consumers see a gap
- · compaction semantics broken
- · “all records for a key read by the same process” — NO LONGER
- · Illustrative numbers — the book’s example says p34.
"When partitioning keys is important, the easiest solution is to create topics with sufficient partitions and NEVER ADD PARTITIONS."
That is about as blunt as this book gets. Combine with Ch. 2's guidance: size on expected future throughput, and remember partition count can only go up, never down — so you must get it right roughly once.
9.5 Custom partitioner — the "Banana" hot-key problem
The scenario: you're a B2B vendor; your biggest customer manufactures handheld devices called Bananas. "You do so much business with customer 'Banana' that over 10% of your daily transactions are with this customer."
With default hash partitioning: Banana's records land in the same partition as other accounts → one partition much larger than the rest → "can cause servers to run out of space, processing to slow down, etc."
What you want: give Banana its own partition, then hash the rest of the accounts across all other partitions.
Banana lands in the same partition as other accounts.
Other records hashed over numPartitions − 1.
public class BananaPartitioner implements Partitioner {
public void configure(Map<String, ?> configs) {}
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes,
Cluster cluster) {
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
int numPartitions = partitions.size();
if ((keyBytes == null) || (!(key instanceOf String)))
throw new InvalidRecordException("We expect all messages " +
"to have customer name as key");
if (((String) key).equals("Banana"))
return numPartitions - 1; // Banana always goes to last partition
// Other records get hashed to the rest of the partitions
return Math.abs(Utils.murmur2(keyBytes)) % (numPartitions - 1);
}
public void close() {}
}Interface: configure, partition, close.
The book's own self-critique, worth internalizing: "Here we only implement partition, although we really should have passed the special customer name through configure instead of hardcoding it in partition." — i.e. hot keys change; recompiling a partitioner to change a customer name is a deploy you shouldn't need.
Note also Utils.murmur2 — this is the same hash Kafka's own default partitioner uses, which is why it's Java-version-stable.