Kafka Partitioning Strategies & Best Practices
When a non-null key is provided, Kafka's DefaultPartitioner uses MurmurHash2 to deterministically map keys to partitions. All records with the same key always land on the same partition β guaranteeing strict ordering per key.
partition = toPositive(murmur2(key)) % numPartitionskey="ACC-001"0x7A4B21C3Partition 2key="ACC-002"0x1D8F405APartition 5key="ACC-001"0x7A4B21C3Partition 2 β samePartitions are the fundamental unit of parallelism, ordering, and throughput in Kafka. Choosing the right partitioning strategy is one of the most consequential design decisions in a Kafka-based system β the wrong choice leads to hot spots, ordering violations, or insufficient parallelism that cannot easily be fixed after data is flowing.
Why Partitioning Matters
The partition key determines:
- Which broker stores the data (leader for that partition)
- Which consumer processes the data (one consumer per partition per group)
- Ordering guarantees (guaranteed within a partition, not across partitions)
The Four Partitioning Strategies
1. Key-Based Partitioning (Default when key is present)
The producer hashes the message key using murmur2 and maps to a partition:
partition = abs(murmur2(keyBytes)) % numPartitions
All messages with the same key are always routed to the same partition β providing ordering guarantees per key.
// Same orderId β always same partition β ordered delivery per order
ProducerRecord<String, OrderEvent> record =
new ProducerRecord<>("orders", orderId, orderEvent);
producer.send(record);
When to use: When ordering per entity matters (events for a given user, order, or device must be processed in sequence).
Pitfall: Key distribution must be uniform. A "hot key" (one key that generates far more messages than others) creates a hot partition β one overloaded broker and one overloaded consumer while others are idle.
2. Round-Robin (Kafka < 2.4, no key)
Keyless messages distributed one-by-one across all partitions. Provides even distribution but poor batching efficiency because each record may go to a different partition before a batch fills.
Kafka β₯ 2.4: Replaced by Sticky Partitioner for keyless messages.
3. Sticky Partitioner (Default when no key, Kafka β₯ 2.4)
The producer sends all keyless messages to the same partition until either batch.size is filled or linger.ms expires, then switches to a new partition. This dramatically improves batching and throughput for keyless producers.
// No key specified β sticky partitioner decides
ProducerRecord<String, ClickEvent> record =
new ProducerRecord<>("clickstream", clickEvent);
When to use: For high-volume event streams where ordering is not required and throughput matters most.
4. Custom Partitioner
Implement org.apache.kafka.clients.producer.Partitioner for business-driven routing:
public class TieredPriorityPartitioner implements Partitioner {
private static final int PREMIUM_PARTITION_COUNT = 2;
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
int totalPartitions = cluster.partitionCountForTopic(topic);
if (keyBytes == null) {
// Null key β uniform distribution over remaining partitions
return ThreadLocalRandom.current().nextInt(PREMIUM_PARTITION_COUNT, totalPartitions);
}
String customerId = new String(keyBytes);
CustomerTier tier = lookupTier(customerId);
return switch (tier) {
case PREMIUM -> Math.abs(customerId.hashCode()) % PREMIUM_PARTITION_COUNT;
case STANDARD -> PREMIUM_PARTITION_COUNT +
Math.abs(customerId.hashCode()) % (totalPartitions - PREMIUM_PARTITION_COUNT);
};
}
@Override
public void close() {}
@Override
public void configure(Map<String, ?> configs) {}
}
partitioner.class=com.example.TieredPriorityPartitioner
Use cases:
- Route premium customers to dedicated partitions with dedicated consumers
- Geographic routing (US β P0-P2, EU β P3-P5)
- Separate low-latency from batch traffic within the same topic
Custom partitioners must be deterministic and stateless β the same key must always map to the same partition, otherwise ordering guarantees break when producer instances restart.
Partition Count Sizing
Getting partition count right is critical. Too few and you can't scale; too many and you add overhead.
The Formula
partitions = max(T/Tp, T/Tc)
Where:
T = target throughput (messages/sec or MB/sec)
Tp = throughput of a single producer partition (measured via perf test)
Tc = throughput of a single consumer partition (your processing speed)
For most production workloads:
| Traffic Level | Recommended Partitions |
|---|---|
| < 10 MB/s | 6β12 |
| 10β100 MB/s | 12β48 |
| > 100 MB/s | 48β200+ |
Sizing Rules of Thumb
- Start with the number of brokers Γ 2 as a baseline (spreads leader partitions evenly)
- Set partition count = expected peak consumers Γ 2 so you have headroom to scale
- Never go below 3 for any production topic (allow for consumer group flexibility)
- Measure, don't guess β use
kafka-producer-perf-test.shto measure your actualTp
# Measure producer throughput per partition (baseline)
kafka-producer-perf-test.sh \
--topic perf-test-1partition \
--num-records 5000000 \
--record-size 1000 \
--throughput -1 \
--producer-props bootstrap.servers=localhost:9092 acks=1
Too Many Partitions: The Overhead
Each partition has costs:
- Per-broker: One open file handle per segment per partition (~3 file handles per partition)
- Per-controller: Metadata for every partition stored in KRaft
- Per-client: Metadata fetch overhead grows with partition count
- Rebalance time: Consumer group rebalances are O(partitions Γ consumers)
Practical upper limit: 4,000β10,000 partitions per broker (depending on broker hardware). Across a 10-broker cluster: 40,000β100,000 total partitions.
Hot Key Problem & Mitigation
What Is a Hot Key?
When a small number of keys generate a disproportionate share of messages, a few partitions receive far more load than others:
Partition distribution with hot key "ad_id = nike_lebron_james":
P0: β 3%
P1: β 2%
P2: β 2%
P3: β 3%
P4: ββββββββββββββββββββββββββββββββββββββββ 90% (HOT - "nike_lebron_james")
Result: The broker hosting the leader replica of P4 experiences severe I/O and network saturation, while the single consumer thread assigned to P4 is overwhelmed, causing consumer lag to climb exponentially. Meanwhile, consumers on P0βP3 sit idle.
Case Study: Ad Click Aggregator Viral Spike
Consider an Ad Click Aggregator system (such as the canonical HelloInterview scenario) where click events are streamed into Kafka to compute billing and real-time CTR (Click-Through Rate).
- The Naive Design: The producer sets
key = ad_id. Under normal conditions, traffic is distributed across all 32 partitions. - The Viral Event: Nike launches a massive LeBron James campaign during the NBA Finals. Tens of millions of users click the ad within minutes.
- The Failure Mode: All clicks for this ad hash to the exact same partition (
abs(murmur2("nike_lebron")) % 32 = 4). Partition 4 is throttled, consumer lag breaches SLA, and memory pressure spikes on broker 4.
The Four Mitigation Strategies
Strategy 1: Omit the Key (Default Sticky Partitioner)β
If strict event ordering is not strictly required (e.g. calculating total click counts or metrics where addition is commutative and associative), simply send records without a key.
Modern Kafka clients use the Sticky Partitioner: records are batched to a single partition until batch.size or linger.ms is reached, then rotates to the next partition in a round-robin cycle.
- Pros: Perfectly balanced traffic across all partitions; maximum batching throughput.
- Cons: Total loss of per-entity FIFO order.
Strategy 2: Random Salting with Two-Stage Aggregationβ
Append a pseudo-random integer suffix (0 to K-1) to the hot key at the producer to distribute the single entity across partitions:
// Producer: Spread hot key across 10 partitions
private static final int SALT_FACTOR = 10;
String partitionKey = isHotKey(adId)
? adId + "#salt=" + ThreadLocalRandom.current().nextInt(SALT_FACTOR)
: adId;
producer.send(new ProducerRecord<>("ad-clicks", partitionKey, clickPayload));
Consumer-Side Two-Stage Aggregation Pattern: In stream processing engines (Kafka Streams or Apache Flink), aggregating a salted key requires two stages:
- Stage 1 (Local Salted Window): Aggregate clicks grouped by
adId#salt=Xover a 1-minute tumbling window. - Stage 2 (Global Merge): Strip the salt suffix and aggregate the intermediate counts by
adIdto compute the final total.
Strategy 3: Compound Key Partitioningβ
Instead of a purely random salt, combine the entity ID with an independent business attribute that naturally distributes load:
adId + "#" + userRegion(e.g.,nike_lebron#us_east,nike_lebron#eu_west)adId + "#" + (userId % 16)
// Compound Key: Ensures all events for a specific user and ad stay ordered,
// while distributing global ad clicks across 16 partition buckets
String compoundKey = adId + "#user_bucket=" + (Math.abs(userId.hashCode()) % 16);
producer.send(new ProducerRecord<>("ad-clicks", compoundKey, clickPayload));
Strategy 4: Producer Backpressure & Dynamic Partition Throttlingβ
In high-scale enterprise topologies, producers or API gateways monitor downstream partition consumer lag:
- If consumer lag on a specific partition exceeds an SLA threshold (e.g. > 50,000 records), the producer dynamically applies backpressure: returning HTTP 429 Too Many Requests to clients or buffering non-critical events in local disk rings until lag subsides.
Ordering Trade-offs
| Scenario | Ordering Guarantee | Approach |
|---|---|---|
| All events for entity X in order | β Per-key | Use entity ID as key |
| Global total ordering | β οΈ Single partition only | 1 partition (no parallelism) |
| No ordering requirement | β None needed | Sticky/round-robin (maximize throughput) |
| Cross-entity ordering | β Not possible | Use a single partition or external coordination |
Ordering with Exactly-Once
For strict ordering with idempotent producers:
// Ordering guarantee: same key = same partition = same sequence
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
// With idempotence enabled, up to 5 in-flight requests maintain ordering
Without idempotence: max.in.flight.requests.per.connection=1 is required for strict ordering with retries (at a significant throughput cost).
Partition Reassignment & Rebalancing
Increasing Partition Count
# Increase from 6 to 12 partitions (irreversible)
kafka-topics.sh --bootstrap-server localhost:9092 \
--alter --topic orders --partitions 12
# WARNING: Existing key-to-partition mappings change!
# "user-123" may now hash to partition 9 instead of partition 3
# Existing messages for "user-123" are still in partition 3
# New messages go to partition 9
# β Ordering broken for in-flight orders
Increasing partitions breaks key-to-partition mapping for existing data. Plan partition counts conservatively upfront or accept a cutover strategy (new topic, migration).
Safe Increase Strategy
# Step 1: Check current assignment
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders
# Step 2: Create reassignment plan for existing data
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--topics-to-move-json-file topics.json \
--broker-list "1,2,3,4" --generate
# Step 3: Execute and monitor
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file reassign.json --execute
# Step 4: Verify completion
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file reassign.json --verify
Best Practices Summary
| Practice | Rationale |
|---|---|
| Use entity ID as key for ordered processing | Same key β same partition β ordered delivery |
| Size partitions at 2Γ expected peak consumers | Leave room to scale without repartitioning |
Use CooperativeStickyAssignor on consumers | Minimize stop-the-world rebalance impact |
| Monitor partition byte rate for hot spots | Detect skew before it causes lag |
| Never use mutable keys | Key hash must be stable or ordering breaks |
| Pre-provision partition count generously | Can increase but can't decrease; increasing breaks key mapping |
Use linger.ms=5β20 for keyless producers | Improves batching efficiency with sticky partitioner |
Interview Questions
Q: What is the difference between key-based and round-robin partitioning?
Key-based partitioning hashes the message key (
murmur2(key) % numPartitions) and routes all messages with the same key to the same partition β guaranteeing ordering per key. Round-robin distributes keyless messages evenly across all partitions on a per-message basis, sacrificing ordering for even distribution. Since Kafka 2.4, keyless messages use the Sticky Partitioner instead of pure round-robin, which improves batching by directing keyless messages to the same partition until a batch fills.
Q: What is the hot key problem and how do you solve it?
A hot key is one key that generates far more messages than average, causing a single partition to receive disproportionate load β overloading one broker and one consumer while others sit idle. Solutions: (1) Key salting β append a random suffix (e.g.,
user-123-0throughuser-123-9) to spread across partitions, then aggregate by original key downstream. (2) Dedicated hot topic β route the hot key to a separate topic with more partitions. (3) Application sharding β divide the hot entity into logical sub-entities with distinct keys.
Q: Why can't you decrease the partition count of a Kafka topic?
Decreasing partitions would require deciding which data to move (or discard) from the removed partitions. More critically, it would invalidate the key-to-partition hash mapping β existing consumers reading historical data from old partitions would find their committed offsets pointing to now-nonexistent partitions, causing data loss or corruption. Kafka's solution is to allow only increases. If you need fewer partitions, create a new topic and migrate.
Q: How do you choose the right partition count?
Use the formula
partitions = max(T/Tp, T/Tc)where T is target throughput, Tp is single-partition producer throughput (measured with perf tests), and Tc is your consumer processing capacity per partition. As a rule of thumb: start at 2Γ your broker count, set it to at least your expected peak consumer count Γ 2, and never go below 3 for production topics. Err on the side of more partitions β it's harder to increase later (breaks key mapping) than to have extra partitions initially.
Related Topics
- Partitions Deep Dive β Partition structure, LEO, HW, and leadership
- Scaling Partitions β Mechanics of increasing partitions and reassignment
- Consumer Groups β How partitions are assigned to consumers
- Kafka Performance Tuning β Throughput and latency optimization
- Exactly-Once Semantics β Ordering guarantees with idempotent producers
Sources
- Apache Kafka Documentation: Partitions
- KIP-480: Sticky Partitioner
- Conduktor Glossary: Kafka Partitioning Strategies
- Narkhede, N., Shapira, G., & Palino, T. β Kafka: The Definitive Guide (O'Reilly)
