When Your Biggest Tenant Breaks Kafka Partition Balance
A single hot tenant can stall your entire Kafka consumer group. Adding partitions won't help — the fix is rethinking your record key. Composite keys (tenant_id:device_id) spread load while preserving per-entity ordering, but you lose global per-tenant order. Here's how to decide and verify.
05 Jul 2025, 02:22 UTC

The hot-tenant problem in plain terms
You run a multi-tenant event pipeline on Kafka. Every event carries a tenant_id and you use that as the record key so each tenant's events stay ordered. Then one tenant grows — maybe a large enterprise customer — and suddenly their partition becomes a bottleneck. The consumer assigned to that partition lags while the other consumers in the group sit idle. Adding more partitions doesn't help because all of that tenant's records still hash to the same single partition.
This is the classic hot-key problem. Kafka guarantees ordering only within a partition, and the default partitioner (murmur2 hash modulo partition count) maps identical keys to the same partition. One hot key means one hot partition, and a consumer group can assign at most one consumer per partition. The fix isn't more partitions — it's a different key strategy.
Why adding partitions doesn't solve a single hot key
Suppose you have a topic with 12 partitions and 100 distinct tenant_id values. The hash spreads them roughly evenly — until one tenant generates 80% of the traffic. That tenant's key hashes to partition 3 (for example). Partition 3 now receives 80% of the bytes. You add 12 more partitions. The murmur2 hash modulo 24 remaps every key, so the hot tenant might move to partition 17 — but it still lands on exactly one partition. The consumer assigned to partition 17 still drowns.
The partition count only increases parallelism for distinct keys. It does not split a single key's stream.
Composite keys: preserve per-entity ordering, spread the load
If your downstream processing only needs ordering per device, per session, or per correlation ID — not per tenant globally — you can change the key to a composite value. A common pattern:
tenant_id + ':' + device_id
The colon is just a delimiter; any character not present in the component values works. The partitioner hashes the full string, so tenantA:device1 and tenantA:device2 land on different partitions. Load spreads across the partition count. Ordering is preserved for each device's event stream, which is often the real requirement (e.g., per-device state machines, per-session funnels).
What you lose: strict global ordering across all devices for that tenant. If your consumer must see every event for tenantA in a single total order, this approach breaks that guarantee. Decide which ordering you actually need.
Worked example: rekeying a tenant-topic
Assume a topic tenant-events with 24 partitions, produced by a Java client on Kafka 3.6 (broker and client). Current key: tenant_id (string). One tenant, acme-corp, produces ~500 MB/s; the other 99 tenants share ~100 MB/s. Partition 7 (where acme-corp hashes) is the hotspot.
- Add a new key field at the producer. Change the serializer to emit
tenant_id + ':' + device_idas the key. Keep the value schema unchanged. - Deploy the producer change. New records fan out across all 24 partitions. The
acme-corptraffic now distributes roughly proportionally to its device cardinality (say 5,000 active devices → ~200 devices per partition on average). - Verify with broker metrics. In Prometheus or JMX, watch
kafka.server:type=BrokerTopicMetrics,name=BytesInPerSecper partition. Before: partition 7 at ~500 MB/s, others at ~4 MB/s. After: all partitions in the 20–30 MB/s range. (Numbers illustrative, not measured.) - Update consumers if they rely on key-based logic. Any consumer that keys its own state store or downstream topic on the original
tenant_idmust now extract the tenant prefix from the composite key.
No topic recreation, no partition count change, no consumer group rebalance beyond the normal cooperative-sticky behavior. The rekey is a producer-side change.
When composite keys aren't enough: key salting
If a single device_id is itself a hot key (e.g., one high-frequency IoT gateway), composite keys still concentrate that device on one partition. Key salting appends a random suffix:
tenant_id + ':' + device_id + '-' + random_int(0, S-1)
S is the salt cardinality (e.g., 10 or 100). The same logical stream now splits across S partitions. Downstream consumers must re-aggregate — typically a two-phase job: a first consumer group reads salted partitions, re-keys by the original logical key, writes to an intermediate topic; a second group consumes the intermediate topic with the original key. This adds latency, operational complexity, and breaks per-key ordering. Reserve salting for analytics pipelines where approximate ordering or windowed aggregation is acceptable. It is not a general fix for stateful per-entity processing.
Null keys and the sticky partitioner
If you ever produce records with null keys, note the behavior shift in the Java client since Kafka 2.4 (KIP-480). Older clients distributed null-keyed records round-robin across partitions. The sticky partitioner now batches null-keyed records to a single partition until the batch fills (controlled by linger.ms and batch.size), then moves to the next partition. This reduces request overhead but can create temporary skew if you produce many null-keyed records in a burst. For keyed workloads, the sticky partitioner behaves like the classic murmur2 partitioner.
Partition count changes remap existing keys
If you do increase the partition count later, the murmur2 hash modulo the new count remaps every existing key. Records for acme-corp:device42 produced before the change live on one partition; records produced after land on a different partition. Consumers reading from earliest offset will see interleaved ordering for that key across the partition boundary. Plan rekeying or migration windows accordingly — either backfill the new key to a new topic, or accept a brief ordering discontinuity during the transition.
Checklist before you rekey
- Confirm the real ordering requirement: per-tenant global, per-device, per-session, or none?
- Measure device/session cardinality for your largest tenants. If cardinality is low (e.g., one device per tenant), composite keys won't spread load.
- Pin the client and broker versions in your documentation. Partitioner defaults differ across Java client, librdkafka, and other libraries.
- Test the new key locally: produce a few thousand records with fixed composite keys to a test topic, then run
kafka-console-consumer --topic test --from-beginning --property print.key=true --property print.partition=trueand verify keys distribute across partitions. - Monitor per-partition
BytesInPerSecand consumer lag after rollout. Expect the hot partition to flatten within one producer batch cycle.
Closing thought
The hot-key problem is a key-design problem, not a partition-count problem. Start by asking what ordering your consumers actually need. If per-device or per-session ordering suffices, a composite key is a low-risk, producer-only change that spreads load immediately. If you need strict global ordering for a single high-throughput entity, you're facing a fundamental throughput ceiling — no key strategy can parallelize a totally ordered stream. In that case, consider whether the entity can be sharded logically (e.g., split the tenant into sub-tenants) or whether the ordering requirement can be relaxed to causal consistency with conflict resolution downstream.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.