Fixing Hot Partitions in Multi-Tenant Kafka
A practical guide to hot partitions: key salting, custom partitioners, and tenant tiering

You've done everything "right." Your Kafka topic is well-provisioned, your consumer group is scaled out, and on paper your throughput should be more than enough. Yet somehow, lag keeps creeping up on the same partition, alerts keep firing for the same consumer, and adding more instances to your consumer group does... nothing.
Sound familiar?
This is the hot partition problem — and it's one of the sneakiest scaling issues in event-driven systems, precisely because your dashboards can look completely fine while it's happening. Aggregate topic throughput? Healthy. Overall consumer group? Green. But dig into per-partition metrics, and you'll find one partition doing 80% of the work while its siblings sit nearly idle.
If you're running Kafka in a multi-tenant system — SaaS platforms, telecom, fintech, anything where one entity (a tenant, a big customer, a power user) can dominate traffic — you will run into this eventually. The good news: it's a well-understood problem with battle-tested fixes.
In this article, I'll break down:
Why hot partitions happen (it's a design trade-off, not a bug)
Three production-proven fixes — key salting, custom partitioners, and tenant tiering — with a deep dive into how a custom partitioner actually works under the hood
An honest take on a "solution" I see floated a lot in production discussions: ditching Kafka consumption for database polling
Let's start with why this happens in the first place.
The Problem: Why Hot Partitions Happen
Kafka distributes messages across partitions using a key:
partition = hash(key) % num_partitions
This gives you a critical guarantee: all messages with the same key land on the same partition, in the order they were produced. That's how Kafka gives you ordering without needing a single global writer/reader.
The catch: Kafka's parallelism unit is the partition, and a consumer group can only assign one consumer per partition. If your key distribution is skewed — say, one tenant generates 80% of your traffic — every one of that tenant's messages hashes to the same partition, no matter how many partitions your topic has.
The result:
That one partition becomes a bottleneck with a hard throughput ceiling (bounded by a single consumer's processing speed).
Every other partition is underutilized.
Consumer lag builds up specifically on that partition — while aggregate topic metrics can look completely healthy, hiding the problem until it's already painful.
If the lagging consumer misses its heartbeat, it can trigger a group rebalance, disrupting the entire consumer group, not just the hot partition.
This isn't a hardware problem. You can't fix it by adding more consumers — Kafka will just leave them idle. It's a key design problem, and it has to be fixed at that level.
Solution 1: Key Salting
The most common and lowest-effort fix: split one hot key into several by appending a suffix.
salted_key = tenant_id + "_" + (bucket_value % N)
Instead of every message for Tenant A hashing to one partition, you now spread it across N sub-keys, which (likely) land on N different partitions.
The critical detail: deterministic vs. random salting.
If
bucket_valueis random, you get pure load spreading — but you lose ordering entirely for that key. Two events from the same entity could land on different partitions and be processed out of sequence.If
bucket_valueis deterministic (e.g.,hash(userId) % N), the same sub-entity always produces the same salted key, and ordering is preserved at that sub-grain.
In practice, this usually means composing the key more specifically to begin with — tenantId + "_" + userId — rather than adding an arbitrary random suffix. If that composite key is still hot (say, one specific user or bot account dominates), that's when a further, deterministic salt on top makes sense — applied surgically to the specific hot sub-key, not blanket across everyone.
Trade-off: you narrow your ordering guarantee to a finer grain than before. That's usually fine — you rarely need ordering across an entire tenant; you need it per user, per order, or per whatever entity is actually stateful.
Solution 2: Custom Partitioner
Salting works by changing the key string and letting Kafka's default hash-based partitioner do its thing. A custom partitioner takes this a step further by giving you explicit control over partition assignment logic, instead of trusting a hash function to spread things out reasonably.
What a Partitioner actually is
When you produce a message, Kafka needs to answer one question: which partition does this message go to? The Partitioner interface is the piece of code that answers that question. Kafka's default partitioner does this:
partition = hash(key) % num_partitions
A custom partitioner means you write your own class implementing this same interface, with your own logic instead of pure hashing, and register it in your producer config.
The interface, broken down
public class HotKeyAwarePartitioner implements Partitioner {
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
...
}
}
This method is called by the producer client every time a message is sent. Kafka passes in everything you might need to decide:
topic— which topic this is going tokey/keyBytes— the message key, as the original object and as raw bytesvalue/valueBytes— the payload (rarely needed for partitioning, but available)cluster— cluster metadata, importantly including how many partitions the topic has
Your job: return an int — the partition this message should go to.
The logic itself
if (isKnownHotKey(key)) {
return dedicatedPartitionFor(key);
}
return defaultHashPartition(keyBytes, cluster);
This is a simple if/else:
isKnownHotKey(key)— a function you write, checking the key against a list you maintain (hardcoded, config-driven, or pulled from a cache/DB) — e.g., is this"tenantA"?If it's a known hot key →
dedicatedPartitionFor(key)picks a partition from a reserved subset you've set aside for it (e.g., partitions 0–4 out of 20).If not → fall back to normal hashing for everyone else, so you're not reinventing logic for the 99% of keys that don't need special treatment.
A concrete filled-in example
Say you have 20 partitions total, and Tenant A (your one hot tenant) should always land in partitions 0–4, while everyone else uses the rest:
public class HotKeyAwarePartitioner implements Partitioner {
private static final Set<String> HOT_TENANTS = Set.of("tenantA");
private static final int HOT_PARTITION_START = 0;
private static final int HOT_PARTITION_COUNT = 5;
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
String tenantId = extractTenantId(key); // e.g. "tenantA_user123" -> "tenantA"
if (HOT_TENANTS.contains(tenantId)) {
// spread this tenant's traffic across its 5 reserved partitions,
// deterministically based on the full key (so same user = same partition)
int bucket = Math.abs(key.hashCode()) % HOT_PARTITION_COUNT;
return HOT_PARTITION_START + bucket;
}
// everyone else: normal hashing across remaining partitions
int numPartitions = cluster.partitionCountForTopic(topic);
return Math.abs(keyBytes.hashCode()) % numPartitions;
}
}
Walking through what happens when a message arrives:
A message comes in with key
"tenantA_user123".extractTenantIdpulls out"tenantA"."tenantA"is inHOT_TENANTS→ we go into the reserved-partition branch.We hash the full key (
"tenantA_user123", not just"tenantA") and mod it by 5 → this decides which one of the 5 reserved partitions this specific user's messages go to.Because the full key hashes deterministically,
"tenantA_user123"always computes to the same bucket → ordering is preserved for that user.A message with key
"tenantB_user456"isn't inHOT_TENANTS, so it falls through to normal hashing across all 20 partitions.
Wiring it in
You register the class in your producer config — Kafka's client library invokes it automatically on every send:
Properties props = new Properties();
props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, "com.yourapp.HotKeyAwarePartitioner");
From that point on, every producer.send() call in your app routes through your custom logic instead of the default one — completely transparent to the rest of your code.
Why this differs from plain salting
With plain salting, you change the key string and let Kafka's default hashing figure out where it lands — you're trusting probability to spread things out reasonably evenly.
With a custom partitioner, you're not trusting probability — you're explicitly deciding the exact partition range a hot key is allowed to use. That's the value-add: predictability and control, at the cost of having to write and maintain this logic yourself instead of just tweaking a key string.
Trade-off: more engineering and operational overhead — you need to detect or maintain a list of hot keys and keep the partitioner logic in sync with actual traffic. Ordering trade-offs are the same as salting; you just get more deliberate control over the blast radius.
Solution 3: Dedicated Topic per Hot Tenant (Tenant Tiering)
If a hot key represents a structurally and permanently high-traffic entity — not a temporary spike — the cleanest fix is architectural: give it its own topic.
topic: tenant-a-events (dedicated, high partition count, dedicated consumer group)
topic: shared-tenant-events (long tail of smaller tenants)
This is the standard pattern in multi-tenant SaaS platforms for isolating "noisy neighbors." You size the dedicated topic's partitions and consumer group specifically for that tenant's load, so it can never starve smaller tenants sharing infrastructure — and within its own dedicated topic, it has room to spread across many partitions or sub-keys without needing to compromise on anything.
Trade-off: more topics to provision, monitor, and reason about, and you need a routing layer at the producer side to decide which topic a given tenant's events go to.
The Database-as-Queue "Solution" — My Honest Take
I want to address this one directly because I see it suggested in production discussions a lot, usually framed as: "just write incoming messages to a DB table, and let consumers poll rows from there instead of consuming from the hot partition."
Technically, it works. With something like Postgres's SELECT ... FOR UPDATE SKIP LOCKED, you can let arbitrarily many consumers pull rows concurrently without contention — you're no longer capped by "1 partition = 1 consumer," so you regain parallelism regardless of key skew.
But here's my honest opinion: I don't think this is a good default fix, and I'd be cautious recommending it.
Here's why:
You're rebuilding a message broker on top of a database. Kafka already gives you ordering guarantees, retention, replay, consumer-group coordination, and backpressure handling — for free, as part of the platform. Move to DB-polling, and you now have to build backpressure, dead-lettering, and offset/cursor tracking yourself.
You haven't removed the hot spot — you've relocated it. A busy queue table under concurrent polling load runs into its own contention issues: lock contention even with
SKIP LOCKED, index bloat, and vacuum pressure in Postgres specifically. The skew doesn't disappear; it just moves from "one Kafka partition" to "one database table under simultaneous read+write load."Polling is inherently worse than push. Kafka consumers use long-polling with near-immediate delivery. A DB-polling consumer either polls aggressively (wasting DB load) or on an interval (adding latency) — there's no free lunch here.
It usually forces you into the outbox pattern anyway, to avoid dual-write inconsistency between your business transaction and your queue table — which is more infrastructure and complexity, not less, layered on top of a system you built specifically to avoid Kafka's complexity.
That said — I'll give credit where it's due: the transactional outbox pattern itself (writing an event to a DB table in the same transaction as your business write, then a separate poller publishes it into Kafka) is a legitimate, widely-used production pattern. But its purpose is solving the dual-write problem for atomicity — not fixing partition skew. If someone's using DB-polling as a permanent replacement for Kafka consumption specifically to dodge hot partitions, I'd treat that as a sign the real fix (salting, custom partitioning, or tenant tiering) wasn't implemented, rather than a valid architectural alternative.
Wrapping Up
The three real fixes — salting, custom partitioning, and tenant tiering — all do the same underlying thing from different angles: they give a hot key more room to spread out, at the cost of narrowing your ordering guarantee to the grain where you actually need it. The database-polling approach doesn't do this at all — it just changes where the bottleneck lives, while adding meaningful complexity back into your system. If you're hitting hot partitions in production, fix the key design first.





