How do you choose the number of partitions for a topic based on target throughput and consumer parallelism?
answer
- partitions = max(T/P, T/C)
- 1 partition → 1 consumer in a group
- consume rate often the binding limit
- can grow partitions, never shrink
- adding partitions reshuffles keys
basics
~20 sPartitions = max(target / per-partition producer rate, target / per-partition consumer rate). Each partition is processed by at most one consumer in a group, so partition count caps consumer parallelism. Round up and leave growth headroom.
solid answer
~40 sPick partitions so you hit your throughput target on both sides and have enough parallelism. Measure (or estimate) sustainable per-partition rates: producer-side P MB/s and consumer-side C MB/s for your processing. For a topic that must handle T MB/s, you need at least ceil(T/P) and ceil(T/C) partitions — take the max. Crucially, in a consumer group each partition is consumed by exactly one member, so partition count is the hard upper bound on consumers that can work in parallel; extra consumers sit idle. Add headroom for growth because you can increase partitions later but never decrease, and increasing them reshuffles the key→partition mapping (breaking per-key ordering across the change). Balance against over-partitioning costs (fds, failover/recovery time, per-broker ceiling). A typical answer: estimate from throughput, sanity-check against per-broker partition limits, round up for growth.
go deeper
Know one partition is consumed by one group member, so partitions cap consumer parallelism.
Compute partitions from target throughput and the slower of produce/consume per-partition rates; round up.
Weigh ordering, key skew, future growth, and per-broker limits when finalizing the count.
Establish partitioning standards balancing parallelism, ordering semantics, and cluster-wide partition budgets.
## What a partition gives you A **partition** is an ordered, append-only log and the unit of both **parallelism** and **ordering** in Kafka. Within a **consumer group** (a set of consumers sharing the load of a topic), Kafka assigns each partition to **exactly one** consumer. So: - Max parallel consumers doing useful work = number of partitions. - Ordering is guaranteed only *within* a partition, never across partitions. ## The throughput sizing formula Let: - `T` = target topic throughput (MB/s) including peak headroom. - `P` = sustainable produce rate per partition (MB/s) on your hardware. - `C` = sustainable consume+process rate per partition (MB/s) — often the binding constraint because downstream processing (DB writes, transforms) is slower than raw I/O. Then: ``` partitions = max( ceil(T / P), ceil(T / C) ) ``` Example: T = 100 MB/s, P = 50 MB/s, C = 10 MB/s (slow processing): produce side needs 2, consume side needs 10 → choose **10**. ## Why round up and over-provision - You can **increase** partitions but **never decrease** them. - Increasing partitions changes the **default hash partitioner** mapping (`hash(key) % numPartitions`), so existing keys can land in different partitions — per-key ordering is not preserved across the change. For keyed topics this is disruptive, so pick a forward-looking count up front. - A common practice: provision for projected 1-2 year peak, not just today. ## Don't over-partition More partitions cost file descriptors, page cache, longer leader failover and crash recovery, and you must stay under the per-broker replica ceiling (~1000s). So the count is a balance: enough for parallelism/throughput and future growth, but bounded by per-broker limits and the overhead of tiny fragmented logs. ## Other inputs - **Ordering requirements**: if you need strict per-key ordering, all records for a key must go to one partition (keyed producing) — partition count then trades off ordering granularity vs parallelism. - **Keyed skew**: a hot key concentrates load on one partition regardless of total count; sizing assumes reasonably even key distribution. - **Rule of thumb cross-check**: keep total cluster partitions within roughly 100 x brokers x RF, and per-broker replicas under your tested ceiling. ## Summary Start from throughput and the slower of produce/consume rates, take the max, round up for growth, and validate against per-broker limits and ordering needs.
- If a consumer group has 12 consumers but the topic has 8 partitions, what happens?Only 8 consumers get a partition and do work; the other 4 sit idle as standby. Partition count caps the effective parallelism of a single group.
- Why is the consumer-side rate often the constraint rather than the producer-side rate?Producing is mostly sequential I/O and very fast, but consumers do real work per record (DB writes, transformations, external calls). That downstream processing is usually slower, so T/C demands more partitions than T/P.
saying these in an interview costs you the question
- Saying more consumers than partitions increases throughput (the extras are idle).
- Sizing only from producer throughput and ignoring slow downstream processing.
- Forgetting partitions can't be reduced and adding them breaks key ordering.
- Assuming ordering holds across partitions.