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.
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.