Lesson 5 / 25

Keys and Partitioning

Choose keys so related events share a partition, and understand how the default partitioner maps keys.

Same key, same partition

When a record has a key, the default partitioner hashes the key bytes with murmur2 and takes the result modulo the number of partitions, so every record with that key goes to the same partition and keeps its order. Records without a key are spread across partitions (sticky batching, then round-robin style) for balance but have no per-entity ordering. Pick a key that groups the events that must stay ordered, for example order_id, customer_id or device_id, while avoiding one very hot key that overloads a single partition.

From record to partition to disk

A producer chooses a partition, batches records, compresses them and waits for acknowledgements.

Four steps: serialise, partition, batch, acknowledge.
Figure 2.1 — Serialise, partition, batch and acknowledge.

Reproducing Kafka's partitioner, run

I implemented murmur2 in Python and compared it with the real broker. For 6 partitions it gives user-1→2, user-2→2, user-3→5, user-4→1, user-5→4, user-6→5, exactly matching the partitions shown by the real consumer earlier.

def murmur2(data: bytes) -> int:
    length = len(data); seed = 0x9747b28c; m = 0x5bd1e995; r = 24
    h = (seed ^ length) & 0xFFFFFFFF
    for i in range(length // 4):
        i4 = i * 4
        k = (data[i4] & 0xff) | ((data[i4+1] & 0xff) << 8) | ((data[i4+2] & 0xff) << 16) | ((data[i4+3] & 0xff) << 24)
        k = (k * m) & 0xFFFFFFFF
        k ^= (k >> r) & 0xFFFFFFFF
        k = (k * m) & 0xFFFFFFFF
        h = (h * m) & 0xFFFFFFFF
        h ^= k
    rem = length % 4; base = length - rem
    if rem == 3: h ^= (data[base+2] & 0xff) << 16
    if rem >= 2: h ^= (data[base+1] & 0xff) << 8
    if rem >= 1:
        h ^= data[base] & 0xff
        h = (h * m) & 0xFFFFFFFF
    h ^= (h >> 13) & 0xFFFFFFFF
    h = (h * m) & 0xFFFFFFFF
    h ^= (h >> 15) & 0xFFFFFFFF
    return h

def partition_for(key: str, n: int) -> int:
    return (murmur2(key.encode()) & 0x7fffffff) % n

for k in ("user-1", "user-2", "user-3", "user-4", "user-5", "user-6"):
    print(k, partition_for(k, 6))

Output:

user-1 2
user-2 2
user-3 5
user-4 1
user-5 4
user-6 5

Key choice is a design decision

Changing a key later reshuffles which partition holds what and breaks per-key ordering across the change. Decide it early, based on the entity whose events need strict order.

Quick check: How does the default partitioner choose a partition for a keyed record?

  • A random number
  • murmur2 hash of the key modulo the partition count
  • The size of the value
  • The broker with most free disk
Answer

murmur2 hash of the key modulo the partition count — Hashing the key makes the choice deterministic: same key, same partition.