How do throughput targets and metrics drive how many partitions a topic should have?
answer
- partitions = max(t/p, t/c)
- 1 partition -> 1 consumer per group (parallelism ceiling)
- can add partitions, never remove
- adding breaks key->partition hashing
- over-partition costs handles, controller, latency
basics
~20 sPartition count must be at least target throughput divided by the per-partition throughput a single producer/consumer can sustain. Since one partition maps to one consumer in a group, partitions also set the max parallelism for consumers.
solid answer
~40 sTwo constraints set partition count. First, parallelism: a partition is consumed by at most one consumer per group, so partition_count is the ceiling on consumer-group parallelism — you need at least as many partitions as your peak consumer instances. Second, throughput: estimate per-partition limits, p (producer MB/s into one partition) and c (consumer MB/s out of one partition), then partitions ≈ max(target/p, target/c). Use measured BytesInPerSec/BytesOutPerSec per partition as inputs. Bias upward for headroom and future growth because increasing partitions later breaks key-based ordering (hashing changes which partition a key lands on). But don't over-partition: each partition costs open file handles, memory, replication connections, and lengthens leader-election/controller work and end-to-end latency. Practical guidance is to keep partitions per broker in the low thousands and total cluster partitions within controller limits.
go deeper
Know partition count caps consumer-group parallelism (one partition per consumer).
Apply partitions = max(t/p, t/c) and know repartitioning is one-way and disruptive.
Balance throughput against over-partitioning costs (handles, controller load, latency, failover).
Define org-wide partitioning standards, growth headroom, and skew/hot-key mitigation across many topics and clusters.
**Why partitions are the unit of scale.** A Kafka topic is split into partitions; each partition is an ordered, independently replicated log. Two things scale with partition count: 1. **Consumer parallelism.** Within a single consumer group, each partition is assigned to exactly one consumer instance. So if a topic has N partitions, at most N consumers in a group do useful work — extra consumers sit idle. Partition count is therefore the hard ceiling on horizontal consumer scaling. 2. **Throughput.** Each partition has practical write and read rate limits set by disk, replication, and the single-threaded-per-partition handling on the consumer side. **The sizing formula** (LinkedIn's classic rule): let `t` be target throughput, `p` the throughput you measure for a single producer to one partition, and `c` the throughput a single consumer sustains from one partition. Then: ``` partitions = max( t/p, t/c ) ``` You get `p` and `c` from measured metrics: BytesInPerSec divided across partitions for producer side, BytesOutPerSec per partition for consumer side, or from benchmarking with kafka-producer-perf-test.sh / kafka-consumer-perf-test.sh. **Round up, with headroom.** Add margin because: - **Repartitioning is disruptive.** Adding partitions changes `hash(key) % partition_count`, so keyed messages start landing on different partitions — ordering-per-key guarantees break for existing keys, and stateful consumers can be thrown off. You can't reduce partitions at all. So pick a number you can live with for the topic's life. - **Future growth** in both data and consumer fleet. **Costs of over-partitioning** (why not just pick 10,000): - **File handles & memory** — each partition replica is open log segments + index files; thousands per broker strain ulimits and page cache. - **Replication overhead** — more leader/follower fetch sessions. - **Controller/metadata load** — leader election and metadata propagation scale with total partition count; with KRaft this is far better than the old ZooKeeper limits, but it's still bounded. - **End-to-end latency** — more partitions can mean smaller batches and more leader elections during failover, lengthening unavailability windows. A broker failure must re-elect leaders for every partition it led. - **Producer memory** — the producer buffers per partition (`batch.size` per partition), so many partitions inflate client memory. **Rules of thumb.** Keep partitions-per-broker in the low thousands; size total cluster partitions to what the controller comfortably handles. A common target is to choose the smallest partition count that meets throughput and parallelism with ~2x headroom. **Edge cases.** - **Key skew.** Even with enough partitions, a hot key concentrates load on one partition; metrics show one partition's BytesIn far above others. Repartitioning won't fix a single dominant key. - **Sticky/cooperative assignment** changes rebalance behavior but not the one-partition-per-consumer ceiling. - **Idempotent/transactional producers** add a small per-partition overhead but don't change the formula.
- Why is increasing partition count later risky?Default partitioning is hash(key) % partitions, so adding partitions remaps keys to different partitions — per-key ordering and any partition-affinity state break. You also can never decrease the count.
- If you add 10 consumers to a group on a 6-partition topic, what happens?Only 6 consumers get a partition each; the other 4 stay idle. Partition count, not consumer count, caps group parallelism.
saying these in an interview costs you the question
- Claiming more consumers than partitions increases throughput.
- Suggesting you can reduce partition count to save resources.
- Ignoring that repartitioning breaks key-to-partition mapping/ordering.
- Recommending huge partition counts without mentioning controller/file-handle/latency costs.