02 / 05

How would you choose a key when all events for a customer must remain ordered?

Difficulty: 4/10
Hot partitions, Partitioning, Ordering

Choosing a Key for Per-Customer Ordering

The key should be the customer identifier, and it must be stable across all events for that customer. If every event for a customer has the same customer ID as its key, the default partitioner will hash that key to the same partition every time, and all events for that customer will be appended to the same partition log. This gives per-customer ordering without any custom partitioning logic. The key must be the same for every event type related to that customer, even if the events are produced by different services. If one service uses customer_id and another uses account_id, the events will land in different partitions and ordering will be lost. This is why a stable, canonical customer identifier is essential, and it should be defined and documented as part of the event contract. The key should also be immutable for the customer; if the customer ID can change, ordering is broken.

The mechanism for this is the default partitioner, which computes murmur2 hash of the key modulo the number of partitions. As long as the number of partitions does not change, the mapping is deterministic. If you add partitions to a topic, the mapping changes for some keys, and events for the same customer may land in different partitions before and after the change, breaking ordering across the partition change. This is why you should choose the partition count carefully up front and avoid changing it unless necessary. If you must add partitions, you need to accept that ordering is preserved only within each partition, and events for a customer may be split across partitions during the transition. A common solution is to use a larger partition count than you need initially, so you have room to grow without changing the count. Another consideration is the key serialization: the key must be serialized consistently across producers. If one producer serializes the customer ID as a string and another as an integer, the bytes differ and the hash differs, so the events land in different partitions. Use a schema registry and a consistent key schema to avoid this.

A common mistake is to use a composite key that includes a timestamp or an event ID, which makes every event unique and breaks per-customer ordering. Another mistake is to use a key that is not stable, such as a session ID or a device ID that changes. A third mistake is to rely on the partitioner to preserve ordering across topics; ordering is per-partition within a topic, and joining across topics does not preserve order unless you design for it. The trade-off is between ordering granularity and parallelism. If you key by customer, you get per-customer ordering but a hot customer can create a hot partition. If you key by a finer-grained entity, such as order ID, you get more parallelism but you lose per-customer ordering. The choice depends on what ordering the business actually requires. Version note: the default partitioner changed in Kafka 2.4 with KIP-480 (sticky partitioner) for null keys, but for non-null keys the murmur2 hash behavior is unchanged. If you rely on the exact partition mapping, be aware of client version differences and test after upgrades.

javascript
  1. 1

    Use a stable customer identifier as the key for all events related to that customer.

  2. 2

    The same key field and serialization must be used by all producers.

  3. 3

    Do not include timestamps or event IDs in the key; that breaks per-customer ordering.

  4. 4

    Adding partitions changes the key-to-partition mapping and can break ordering across the change.

  5. 5

    Choose a larger partition count than you need to leave room for growth.

  6. 6

    Avoid unstable keys like session IDs or device IDs that change.

  7. 7

    The default partitioner is murmur2 hash of the key; behavior is stable for non-null keys.

Share

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