Skip to main content

Kafka Partitioning Strategies & Best Practices

Kafka Partitioning Strategies
Guarantees per-key ordering

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)) % numPartitions
Assignment Examples
key="ACC-001"0x7A4B21C3Partition 2
6 partitions
key="ACC-002"0x1D8F405APartition 5
6 partitions
key="ACC-001"0x7A4B21C3Partition 2 βœ“ same
6 partitions
Production Gotchas
⚑Hot partition risk: if one key dominates traffic (e.g. a single user_id), one partition gets overloaded
⚑Adding partitions breaks existing key β†’ partition mapping (rebalance required)
⚑Key null β†’ falls back to Sticky Partitioner (since Kafka 2.4)
⚑Use custom partitioner to override: sameKey β†’ samePartition, but control hot-spot routing

Partitions 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:

  1. Which broker stores the data (leader for that partition)
  2. Which consumer processes the data (one consumer per partition per group)
  3. 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
warning

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 LevelRecommended Partitions
< 10 MB/s6–12
10–100 MB/s12–48
> 100 MB/s48–200+

Sizing Rules of Thumb

  1. Start with the number of brokers Γ— 2 as a baseline (spreads leader partitions evenly)
  2. Set partition count = expected peak consumers Γ— 2 so you have headroom to scale
  3. Never go below 3 for any production topic (allow for consumer group flexibility)
  4. Measure, don't guess β€” use kafka-producer-perf-test.sh to measure your actual Tp
# 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).

  1. The Naive Design: The producer sets key = ad_id. Under normal conditions, traffic is distributed across all 32 partitions.
  2. The Viral Event: Nike launches a massive LeBron James campaign during the NBA Finals. Tens of millions of users click the ad within minutes.
  3. 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 KK 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:

  1. Stage 1 (Local Salted Window): Aggregate clicks grouped by adId#salt=X over a 1-minute tumbling window.
  2. Stage 2 (Global Merge): Strip the salt suffix and aggregate the intermediate counts by adId to 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

ScenarioOrdering GuaranteeApproach
All events for entity X in orderβœ… Per-keyUse entity ID as key
Global total ordering⚠️ Single partition only1 partition (no parallelism)
No ordering requirement❌ None neededSticky/round-robin (maximize throughput)
Cross-entity ordering❌ Not possibleUse 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
caution

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

PracticeRationale
Use entity ID as key for ordered processingSame key β†’ same partition β†’ ordered delivery
Size partitions at 2Γ— expected peak consumersLeave room to scale without repartitioning
Use CooperativeStickyAssignor on consumersMinimize stop-the-world rebalance impact
Monitor partition byte rate for hot spotsDetect skew before it causes lag
Never use mutable keysKey hash must be stable or ordering breaks
Pre-provision partition count generouslyCan increase but can't decrease; increasing breaks key mapping
Use linger.ms=5–20 for keyless producersImproves 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-0 through user-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.


Sources

  1. Apache Kafka Documentation: Partitions
  2. KIP-480: Sticky Partitioner
  3. Conduktor Glossary: Kafka Partitioning Strategies
  4. Narkhede, N., Shapira, G., & Palino, T. β€” Kafka: The Definitive Guide (O'Reilly)
πŸ“–
Track Page Progress0 / 635 Read
Knowledge Base Completion0%