Key by the finest ordering unit, give each tenant a configurable spread across partitions, and isolate the largest tenants entirely
I start by asking what ordering is actually needed. Almost never is it per tenant. It is per entity inside a tenant (per user, order, device, or aggregate). That distinction is what frees the design. If I key by tenantId alone, a large tenant is pinned to one partition and one consumer, which is the hot partition problem. If I key by the entity within the tenant, load spreads, per-entity ordering is preserved, and tenant-level ordering is given up, which the business usually does not need. So the baseline is a composite key such as tenantId:entityId, with the partitioner designed so a tenant's traffic lands on a controlled set of partitions.
Tenants are not uniform, so I use tiers. Small tenants (the long tail) share a pool topic and hash to one partition by tenant, which keeps their footprint small and preserves tenant-level ordering if someone cares. Medium tenants are spread across a bounded number of partitions, a bucket count B per tenant, by hashing the entity ID inside that range. Very large tenants get a dedicated topic (or a dedicated cluster for regulated or extreme cases), which gives isolation: their own partitions, retention, quotas and consumer group, so their spikes and lag cannot affect others. Tenants are promoted between tiers by measured throughput, not by sales promises. Capacity planning then becomes per tier: partitions are sized for the tier's aggregate peak, and quotas per client or user stop one tenant from monopolizing broker bandwidth.
Ordering trade-off: per-entity ordering is kept, tenant-wide ordering is not for tenants with buckets above 1. Make that contract explicit in the tenant-event documentation.
Changing a tenant's bucket count remaps its entities, so ordering can break across the change boundary. Treat bucket changes as a controlled migration: version the config, drain or fence in-flight events, and change during low traffic, or move the tenant to a new topic.
All producers for the topic must share the same partitioner logic and config version. A mismatch silently breaks ordering and locality, so ship the partitioner as a shared library with contract tests.
Isolation: dedicated topics for large tenants cost more partitions and operations, but give independent retention, quotas, lag alerting and blast radius. Pick the threshold from measured share of traffic and SLA, not from a fixed number.
Fairness and capacity: apply produce and fetch quotas per principal or client-id so one tenant cannot saturate broker network. Monitor per-tenant bytes in, per-partition skew and consumer lag by tenant.
Consumer side: large tenants need enough consumers per partition window. Consider a key-aware worker pool inside a consumer if per-partition processing is the bottleneck.
Alternative: one topic per tenant. It gives the cleanest isolation and per-tenant retention and ACLs, but breaks down with thousands of tenants because of partition and metadata overhead, so I use it only for the top tier.
Common mistake: adding partitions to fix a hot tenant when the key is tenantId. The tenant still hashes to one partition. Another is salting randomly with no entity affinity, which destroys per-entity ordering unnecessarily.
Version note: KRaft (the only mode from Kafka 4.0) raises practical partition limits compared with ZooKeeper-era guidance, but dedicated-topic-per-tenant designs should still be load-tested on your version and hardware.
0-2 years experience
2-5 years experience
5-8 years experience
8+ years experience