Learn Labs
3. Kafka Producers: Writing Messages to Kafka

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 ProducerRecord with just a topic and a value, with the key set to null by default, most applications produce records with keys."

Keys serve two goals:

  1. Additional information stored with the message.
  2. 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-robin2.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.)

KeyMaps overConsequence
Null keyavailable partitionsnever blocked by an offline partition
Non-null keyall partitionscan 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 UniformStickyPartitioner will 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."*

BEFORE — 10 partitionsuser 045189 → hash % 10 = p34 · history lives in p34AFTER adding partitions — 16user 045189 → hash % 16 = p7 · new records go to p7a single key’s history is now SPLITacross two partitions
  • · 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.
Figure 3.9.39.4 ⚠️ Adding partitions breaks key→partition mapping

"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.

DEFAULT HASH — hot partition
p0
others
p1
others + BANANA
p2
others
p3
others

Banana lands in the same partition as other accounts.

CUSTOM — isolated hot key
p0
others
p1
others
p2
others
p3
BANANA — dedicated

Other records hashed over numPartitions − 1.

Figure 3.9.49.5 Custom partitioner — the 'Banana' hot-key problem
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.


On this page