Lesson 13 / 25

Choosing the Partition Count

Size partitions from target throughput and consumer parallelism and know the cost of changing later.

A sizing rule of thumb

Estimate the target throughput (MB/s), measure what a single partition can handle for producers and for your consumer logic, and take the larger of target / producer_per_partition and target / consumer_per_partition. Add headroom for growth, because you can increase partitions later but never decrease them, and increasing changes the key-to-partition mapping (new records with an existing key may land in a different partition, breaking per-key order across the change). Too many partitions have costs too: more open files, longer leader elections and more memory. Hundreds per topic are common; tens of thousands per cluster need care.

Choose partitions, keys and lifetime

How many partitions, which key and how long to keep data are the three decisions that shape a topic.

Three decisions: partitions, key, retention.
Figure 4.1 — Partitions, key and retention.

Sizing arithmetic, run

I ran this: a 300 MB/s target with 50 MB/s per partition for producers and 30 MB/s per partition for consumers needs max(6, 10) = 10 partitions.

import math
target_mb, prod_mb, cons_mb = 300, 50, 30
print(max(math.ceil(target_mb / prod_mb), math.ceil(target_mb / cons_mb)))

Output:

10

Adding partitions changes key placement, run on a real broker

I raised orders from 6 to 8 partitions and produced the six keys again. Using the murmur2 function from earlier (repeated here so the snippet runs on its own), the real broker placed them exactly as the 8-partition mapping predicts: user-1 moved from partition 2 to partition 4, user-2 from 2 to 0, user-3 from 5 to 3.

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

print({k: partition_for(k, 8) for k in ("user-1", "user-2", "user-3", "user-4", "user-5", "user-6")})
# kafka-topics.sh --alter --topic orders --partitions 8   (then produce the same keys again)

Output:

{'user-1': 4, 'user-2': 0, 'user-3': 3, 'user-4': 7, 'user-5': 6, 'user-6': 5}

Quick check: What is true about changing a topic's partition count?

  • You can freely decrease it
  • You can increase it but not decrease it, and key placement may change
  • It never affects ordering
  • It is fixed forever
Answer

You can increase it but not decrease it, and key placement may change — Partitions can only grow, and a different modulus changes where existing keys map.