04 / 05

How would you scale a topic dominated by a small number of keys?

Difficulty: 8/10
Large-scale design, Partitioning, Capacity planning

Scaling a Topic Dominated by a Small Number of Keys

When a topic is dominated by a small number of keys, the problem is hot partitions: the keys that dominate map to a few partitions, and those partitions become the bottleneck. Adding partitions or consumers does not help because the hot key still maps to one partition, and all records for that key must go to the same partition to preserve ordering. There are three main strategies: key redesign, key sharding, and a hybrid approach. Key redesign means changing the key to a higher-cardinality field that still captures the ordering requirement. For example, if the key is customer ID and one customer dominates, and the ordering requirement is per order, change the key to order ID. This spreads the load across more partitions but changes the ordering guarantee. Key sharding means splitting the hot key into multiple sub-keys, such as customer ID + (order ID % 10), which spreads the load across 10 partitions per customer. This preserves ordering within each sub-key but breaks ordering across sub-keys. The hybrid approach uses a shared topic for the long tail and dedicated topics for the largest keys, with a routing layer that decides which topic a record goes to. Each strategy has different trade-offs between ordering, balance, and complexity.

The mechanism for key sharding is to modify the key before producing, usually by appending a shard suffix derived from a field in the record. For example, key = customer_id + "-" + (order_id.hashCode() % 10). This changes the key that the partitioner hashes, so records for the same customer go to different partitions depending on the shard. The consumer must be aware of the sharding and must aggregate across shards if it needs per-customer state. This is a significant change: it moves the complexity from the broker (which cannot split a partition) to the application (which must handle partial ordering). The mechanism for the hybrid approach is to maintain a routing table that maps keys to topics. Large keys get their own topic (or a set of partitions within a dedicated topic), and small keys share a topic. The routing table must be consistent and highly available, and it must be updated when a key grows or shrinks. The trade-off is between ordering and scalability. Key redesign sacrifices some ordering; key sharding sacrifices ordering within the hot key; the hybrid approach preserves ordering but adds routing complexity. Version note: Kafka does not provide built-in hot-key detection or sharding; you must implement it in the producer or in a stream processor. Some managed services provide hot-key detection, but the responsibility for the design remains yours.

A common mistake is to add partitions without changing the key, which does not help because the hot key still maps to one partition. Another mistake is to shard the key without understanding the ordering implications, which can cause incorrect results if the consumer assumes per-key ordering. A third mistake is to use a random key for the hot entity, which destroys ordering completely. The trade-off is between the business requirement for ordering and the technical need for scalability. If the business can tolerate per-order ordering instead of per-customer ordering, key redesign is the simplest solution. If the business requires per-customer ordering even for the hot customer, the hybrid approach is the only option that preserves it. If the business can tolerate partial ordering (e.g., ordering within a region or a time window), key sharding is a middle ground. Version note: the choice of shard count is a trade-off: more shards give more parallelism but more complexity and more partial ordering. A common approach is to use a small number of shards (e.g., 10 or 16) and to adjust based on the load. Also, if you use Kafka Streams, you can use the Processor API to implement custom sharding and aggregation.

javascript
  1. 1

    Hot keys map to a few partitions; adding partitions does not help.

  2. 2

    Key redesign: change to a higher-cardinality key, sacrificing some ordering.

  3. 3

    Key sharding: split the hot key into sub-keys, preserving per-sub-key ordering.

  4. 4

    Hybrid: dedicated topics for large keys, shared topic for the long tail.

  5. 5

    Sharding requires the consumer to handle partial ordering and aggregation.

  6. 6

    The hybrid approach preserves per-key ordering but adds routing complexity.

  7. 7

    Monitor partition rates to detect hot partitions.

  8. 8

    Kafka does not provide built-in hot-key detection or sharding.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.