skip to content

How do you choose the number of partitions for a topic based on target throughput and consumer parallelism?

level: middleimportance: should knowfreq 55%

answer

  1. partitions = max(T/P, T/C)
  2. 1 partition → 1 consumer in a group
  3. consume rate often the binding limit
  4. can grow partitions, never shrink
  5. adding partitions reshuffles keys

basics

~20 s

Partitions = 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 s

Pick 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

for a junior

Know one partition is consumed by one group member, so partitions cap consumer parallelism.

for a middle

Compute partitions from target throughput and the slower of produce/consume per-partition rates; round up.

for a senior

Weigh ordering, key skew, future growth, and per-broker limits when finalizing the count.

for a principal

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.

context