A Kafka topic is split into partitions, and the partition is the unit everything else is built on: ordering is guaranteed only within a partition, each partition is read by at most one consumer in a consumer group, and each partition has one leader broker that takes all of its writes. Which partition a record lands in is decided by the producer, before the record leaves the client. That decision fixes which broker does the work, which consumer sees the record and which other records it is ordered with.
How partitions are stored, replicated, sized and reassigned on the brokers is covered in Kafka partition architecture. This page covers the producer's side: the rules the Java client uses to choose a partition, how those rules changed across releases, why two clients can disagree about where a key belongs, how records are batched per partition, and how to find and fix a hot partition. Configuration names and defaults were checked against the Apache Kafka 4.1 producer configuration documentation on 2026-10-02.
The partition choice, in order
The Java producer decides the partition for each record by the first rule that applies. If the record names a partition explicitly, that partition is used. Otherwise, if partitioner.class is set, the custom partitioner decides. Otherwise the built-in logic runs, and the documentation describes it in two cases: if a key is present, choose a partition based on a hash of the key; if no key is present, choose the sticky partition, which changes when at least batch.size bytes have been produced to it.
For keyed records the hash is murmur2 over the serialized key bytes, made positive and taken modulo the number of partitions. Two consequences follow directly. The partition depends on the serialized bytes, so changing the key serializer, or adding a field to a JSON key, moves every key. And it depends on the partition count, so adding partitions to a topic moves most keys to new partitions, which breaks per-key ordering across the change and splits any state keyed by partition. Plan the count up front, as the architecture article describes, rather than relying on growing it later.
import org.apache.kafka.common.utils.Utils;
import java.nio.charset.StandardCharsets;
// The Java client's key-hash rule: positive murmur2 of the serialized key, modulo partition count.
static int partitionForKey(byte[] serializedKey, int numPartitions) {
return Utils.toPositive(Utils.murmur2(serializedKey)) % numPartitions;
}
// Same answer the producer computes, provided the key serializer matches (here: StringSerializer).
int p = partitionForKey("tenant-42".getBytes(StandardCharsets.UTF_8), 12);Two producer settings modify the built-in logic. partitioner.ignore.keys (default false) makes the producer ignore keys when choosing a partition, so keyed records are spread like unkeyed ones; keys are still stored and still used by log compaction. partitioner.adaptive.partitioning.enable (default true) lets the producer send more records to partitions on faster brokers. Adaptive behaviour only affects records whose partition is not fixed by a key hash: unkeyed records, or all records when keys are ignored. A keyed record always goes where its hash says, however slow that partition's leader is. Both settings have no effect when a custom partitioner is configured. A third, partitioner.availability.timeout.ms (default 0, meaning off), lets the partitioner treat a partition as unavailable when its broker has not been able to accept produce requests for that long.
How unkeyed records used to be spread, and why it changed
Older clients sent unkeyed records round-robin, one record per partition in turn, which spread load evenly but produced tiny batches and latency that rose with partition count. KIP-480, released in Kafka 2.4, introduced the sticky partitioner: stick to one partition until its batch is sent, then pick another. Batches got bigger and latency fell.
KIP-794 then showed that the sticky partitioner was not actually uniform. Stickiness lasted until a batch was drained, and slower brokers drain batches more slowly, so a slow broker held the sticky partition longer and received more records, which made it slower still. KIP-794, accepted in 2022 and shipped in the 3.3 line, switched to the rule the documentation states today, switching partitions after a fixed number of bytes, and added the adaptive option that steers unkeyed traffic away from slow brokers. The classes DefaultPartitioner and UniformStickyPartitioner were deprecated; the replacement is to leave partitioner.class unset and use partitioner.ignore.keys where you want key-agnostic spreading. If old configuration still names those classes, remove it.
Two clients, two answers for the same key
The murmur2 rule is a property of the Java client, not of Kafka. Clients built on librdkafka, including confluent-kafka for Python, Go and .NET, default to a partitioner called consistent_random, which hashes keys with CRC32. A Java service and a Python service producing the same key to the same topic will, by default, send it to different partitions. Nothing fails; per-key ordering silently disappears, and joins that assume co-partitioning return wrong answers rather than errors.
librdkafka offers murmur2 and murmur2_random as Java-compatible options; murmur2_random is described as functionally equivalent to the Java producer's default, hashing keys with murmur2 and spreading NULL keys randomly. Set it explicitly, for example "partitioner": "murmur2_random" in a confluent-kafka configuration, in every non-Java producer that shares keyed topics with Java producers.
Batching happens per partition
The producer's accumulator keeps one open batch per partition. A sender thread ships a partition's batch when it reaches batch.size bytes (default 16,384) or when linger.ms has elapsed since the batch was opened (default 5 in the 4.x documentation; older clients defaulted lower, so check your version). All open batches share buffer.memory, 32 MiB by default, and send() blocks when it is full. Compression is applied per batch, so small batches also compress worse.
This makes partition count a producer cost, not only a broker one. A producer writing keyed records spread over 600 partitions has up to 600 batches filling in parallel, each receiving a six-hundredth of the traffic. At modest throughput most batches are sent by the linger timer half-empty, requests multiply, compression ratios fall and broker request rates rise. Raising linger.ms trades latency for fuller batches; raising batch.size helps only when traffic per partition can actually fill the larger batch. When latency budgets are tight, fewer partitions per producer is the more effective fix.
Hot keys: measuring skew
Key hashing spreads distinct keys evenly, but it cannot spread traffic that is concentrated in a few keys. Real keys are rarely uniform: tenants, devices and users follow heavy-tailed distributions, and one key with 15 percent of the traffic will put at least that much on one partition however many partitions you add. The simulation below draws 200,000 records from 5,000 tenant keys with a Zipf-like distribution and assigns them to 12 partitions. It uses an MD5-based stand-in hash, so the specific partitions differ from what murmur2 would choose, but the skew it shows is the property of the key distribution, not of the hash.
import hashlib
import random
from collections import Counter
PARTITIONS = 12
def partition_for(key: str, n: int = PARTITIONS) -> int:
# Stand-in hash for illustration only; the Java client uses murmur2.
return int.from_bytes(hashlib.md5(key.encode()).digest()[:4], "big") % n
def zipf_keys(n_keys, n_records, s=1.1, seed=3):
rng = random.Random(seed)
weights = [1 / (r ** s) for r in range(1, n_keys + 1)]
keys = [f"tenant-{r}" for r in range(1, n_keys + 1)]
return rng.choices(keys, weights=weights, k=n_records)
def skew(load):
counts = [load.get(p, 0) for p in range(PARTITIONS)]
return max(counts) / (sum(counts) / PARTITIONS)
records = zipf_keys(5_000, 200_000)
plain = Counter(partition_for(k) for k in records)
print(f"top key share = {Counter(records).most_common(1)[0][1] / len(records):.3f}")
print(f"max/mean, by key = {skew(plain):.2f}")
hot = {k for k, _ in Counter(records).most_common(3)}
rng = random.Random(9)
salted = Counter(
partition_for(f"{k}#{rng.randrange(8)}") if k in hot else partition_for(k)
for k in records
)
print(f"max/mean, salted = {skew(salted):.2f}")Running it prints:
top key share = 0.157
max/mean, by key = 2.32
max/mean, salted = 1.92One tenant carries 15.7 percent of all records, almost double the 8.3 percent an even share of 12 partitions would be, so the busiest partition receives 2.32 times the mean. Spreading the three hottest keys over eight sub-keys each brings it to 1.92, not to 1.0, because the remaining 4,997 keys still land unevenly and hashed sub-keys can collide. Salting reduces skew; it does not remove it.
In production, measure skew from the broker side rather than guessing: per-partition bytes-in and messages-in rates from broker metrics, per-partition log size growth, and per-partition consumer lag. A partition whose lag grows while its siblings stay flat is almost always a hot key, not a slow consumer. Lag and the consumer side of this are covered in Kafka consumers in depth.
Fixing a hot key
Every fix trades something, usually per-key order. Pick by what the consumers of the topic need.
| Fix | How | What you give up |
|---|---|---|
| Salt the hot key | Append a sub-key, such as a random or round-robin suffix from 0 to N-1, for known hot keys | Per-key order across sub-keys; consumers must merge or tolerate it |
| Finer key | Key by tenant plus entity id instead of tenant | Order only per entity, which is often all you needed |
| Dedicated topic | Route the hot tenant to its own topic with its own partitions | More topics to operate; routing logic in producers |
| Two-stage aggregation | Pre-aggregate with a salted key, then re-key by the real key for the final result | An extra hop and extra latency |
| Ignore keys | partitioner.ignore.keys=true on a topic that does not need key order | All per-key order on that topic |
The custom partitioner below applies salting inside the producer, so application code keeps writing the natural key. Hot keys are configured, hashed to a home partition and then spread over the next few partitions. Because the record keeps its original key, log compaction still sees one key; but compaction keeps only the latest record per key per partition, so a salted key leaves one latest record in each of several partitions, which a consumer rebuilding state must reconcile. Use salting only on topics where that is acceptable, and see Kafka ordering in depth for the order guarantees it removes.
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ThreadLocalRandom;
import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.utils.Utils;
/** Spreads a configured set of hot keys over SPREAD partitions; everything else hashes normally. */
public class HotKeySpreadingPartitioner implements Partitioner {
private static final int SPREAD = 4;
private Set<String> hotKeys = Set.of();
@Override public void configure(Map<String, ?> configs) {
Object v = configs.get("hot.keys"); // e.g. "tenant-1,tenant-2"
if (v != null) hotKeys = Set.of(v.toString().split(","));
}
@Override public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
int n = cluster.partitionCountForTopic(topic);
if (keyBytes == null) { // custom partitioners must handle null keys
return ThreadLocalRandom.current().nextInt(n);
}
int home = Utils.toPositive(Utils.murmur2(keyBytes)) % n;
if (key != null && hotKeys.contains(key.toString())) {
return (home + ThreadLocalRandom.current().nextInt(SPREAD)) % n; // per-key order is given up
}
return home;
}
@Override public void close() { }
}
Partition count from the client's side
Partitions cap consumer parallelism: a group with more consumers than partitions leaves the extras idle, so the count must cover the peak number of consumers you will want. Each partition also costs batching efficiency in producers and replication work on brokers. When consumers are added or removed, partitions move between them; consumer group rebalancing explains the protocols that do that and their costs.
Worked example: a multi-tenant event topic
A SaaS platform writes audit events to a 24-partition topic, keyed by tenant id, from Java API servers and a Python ingestion worker. Consumers build per-tenant timelines and need events for a tenant in order. Two symptoms appear: one partition's consumer lags by hours during business hours, and some tenants' timelines show events out of order.
The lag is a hot tenant: the top tenant produces 22 percent of events, and broker metrics show its partition at more than five times the mean bytes-in, as 22 percent of 24 partitions' traffic must be. The ordering problem is the Python worker, which uses the librdkafka default partitioner, so the same tenant lands in a different partition from the Java servers. The team sets partitioner to murmur2_random on the Python worker, after which each tenant's events share one partition again. For the hot tenant they change the key to tenant plus resource id, because timelines only need order per resource, which spreads that tenant over many partitions and keeps order where it matters. They do not add partitions, since that would have remapped every tenant during the migration.
Failure modes
- Mixed-client hashing: Java and librdkafka producers on one topic with default partitioners put the same key in different partitions.
- Key format drift: a serializer or key schema change moves every key, silently breaking order and co-partitioned joins.
- Growing the count: adding partitions remaps most keys; stateful consumers see keys appear in new partitions mid-stream.
- Hot partition: one heavy key overloads one leader and one consumer while others idle; adding partitions does not help.
- Half-empty batches: keyed traffic spread over many partitions sends small batches, raising request rates and lowering compression.
- Stale partitioner config: deprecated partitioner classes kept in configuration after upgrades.
What to do next
- List every producer of each keyed topic and its client library; set
partitioner=murmur2_randomon every librdkafka-based producer that shares keys with Java producers. - Remove
DefaultPartitionerandUniformStickyPartitionerfrom configuration; usepartitioner.ignore.keysif you need key-agnostic spreading. - Chart per-partition bytes-in and consumer lag for your busiest topics and flag any partition above twice the mean.
- For each hot key, choose a fix from the table based on the order your consumers actually need.
- Check average batch size and request rate per producer; if batches are small, reduce partitions per producer or raise
linger.mswithin your latency budget. - Write down the key serializer and partition count as part of each topic's contract, and treat changes to either as a migration.