Sizing Partitions

Sizing Partitions and Replication

Partitions are Kafka 129 's unit of parallelism, but they are not free. A common starting point is partitions = max(T / P, T / C), where T is the target throughput and P and C what one producer and one consumer can handle per partition, rounded up for growth, since adding partitions later moves keys to new partitions (Topics and Partitions). Each partition costs file handles and memory on its brokers, a leader to elect on failover, and, on the client side, smaller batches: a keyed producer fills one batch per partition, so the same stream spread over more partitions leaves in more, smaller requests. l0615_partitions.py sent 200,000 BookNest events keyed by order_id (acks=all, zstd 126 , linger.ms=5) to topics with 1, 3, 12 and 48 partitions on the one broker, three rounds:

Output of 88
...
round 2 partitions=1    223,767 msg/s  2326/batch p50   19.2 ms  p99   50.7 ms  0 failed
round 2 partitions=3    200,026 msg/s   893/batch p50   19.5 ms  p99  161.4 ms  0 failed
round 2 partitions=12   182,905 msg/s   237/batch p50  112.6 ms  p99  265.9 ms  0 failed
round 2 partitions=48    66,992 msg/s    73/batch p50 1460.7 ms  p99 2028.8 ms  0 failed
...
booknest.p1: 1 directories, 5 files
booknest.p48: 48 directories, 240 files

Across the three rounds, one partition carried 3.3-8.1 times what 48 did (224,000-309,000 against 38,000-67,000 events a second), batches shrank from 2,100-3,400 records to 57-220, and every partition added five files (the segment, its offset and time indexes, a leader-epoch checkpoint and partition metadata). On one broker, extra partitions buy nothing. On a cluster, choose enough to spread load over every broker and run the consumers you need, not ten times that: BookNest's order topic keeps 3, and 12 would be generous at 100 times today's traffic.

For durability, keep The In-Sync Replica Set's rule: replication factor 3, min.insync.replicas=2, acks=all, which survives one broker failure without losing acknowledged writes. It triples disk use, and in the cloud the copies crossing zones are a cost line of their own (MSK, Event Hubs, Google Kafka).