Command Palette

Search for a command to run...

Hectal
PHASE 1Beginner ~9 min· topic 4 of 5

Topic 1.4

Choosing the Partition Count

In one line

Partition count sets maximum consumer parallelism and spreads load, but each partition costs broker memory, file handles, replication work, metadata, failover and rebalance time. Size it from target throughput divided by measured per-partition throughput for producers and consumers, plus headroom, and plan for growth because increasing it later remaps keys.

0/5 · 0%

Think of it like this

Opening checkout lanes in a new supermarket. Too few and queues grow on busy days; too many and you pay staff to stand idle, and it's hard to remove lanes once customers are used to them.

Key ideas

  1. 01

    Throughput formula: partitions ≥ max(target MB/s ÷ producer MB/s per partition, target MB/s ÷ consumer MB/s per partition). The consumer side usually dominates, because consumers do real work (database writes, API calls).

  2. 02

    Measure per-partition throughput with your own record sizes, compression and processing: kafka-producer-perf-test.sh and kafka-consumer-perf-test.sh for raw numbers, and a load test of your actual consumer for processing throughput. Don't rely on universal numbers.

  3. 03

    Consumer parallelism: a group can't use more consumers than partitions. If you expect to scale to 40 consumer instances, you need at least 40 partitions.

  4. 04

    Costs of many partitions: more open files and memory on brokers, more replication fetch requests, longer leader elections and recovery after broker failure, more metadata for clients, slower rebalances, and less efficient batching (records spread thinner). KRaft raised practical cluster-wide limits a lot, but per-broker limits still matter (commonly a few thousand partitions per broker).

  5. 05

    Growth: adding partitions later changes hash(key) % N, so the same key starts landing in a different partition, breaking per-key ordering during the transition. Choose a count with headroom (for example 2× expected need), or plan migration to a new topic.

  6. 06

    Keyed topics: choose a count that divides cleanly among expected consumer counts (12, 24, 48 are popular because they divide by 2, 3, 4, 6), so assignments are even.

Code & diagrams

perf-test.shbash
# Producer throughput for 1 KB records, zstd, acks=all, into a 1-partition test topic
kafka-producer-perf-test.sh --topic perf-1p --num-records 2000000 --record-size 1024 \
  --throughput -1 --producer-props bootstrap.servers=$B acks=all compression.type=zstd linger.ms=10 batch.size=131072
2000000 records sent, 61237.8 records/sec (59.80 MB/sec), 38.1 ms avg latency, 212.0 ms max latency

# Consumer throughput
kafka-consumer-perf-test.sh --bootstrap-server $B --topic perf-1p --messages 2000000
start.time, end.time, data.consumed.in.MB, MB.sec, data.consumed.in.nMsg, nMsg.sec
..., 1953.1, 142.4, 2000000, 145838.9
# Numbers vary hugely by hardware, record size and settings. Measure yours.

Interview problem

The problem

1 million events/sec, 20 consumers, 3 brokers

You need 1M events/sec, you'll run 20 consumer instances, and you have 3 brokers. How do you reason about partition count? Don't give a magic number; explain the capacity-planning process.

You're given

  • 1M events/sec
  • Average event 500 bytes
  • 20 consumers today
  • 3 brokers

The interviewer follows up

01

Why not just create 5,000 partitions to be safe?

When it breaks

Partition count increased on a keyed topic in production

What you see

Keys remap; new events for existing orders land in different partitions than old ones. Consumers processing per-key state see events out of order during the transition.

Fix & prevent

Size with headroom up front; if you must grow, pause producers until consumers drain, or migrate to a new topic with a controlled cutover.

Explain it without notes

01

List five costs of having too many partitions.

Practice

01

Run the producer and consumer perf tests in your lab with 1 KB records and compute partitions needed for 50 MB/s.

Trade-offs

  • ↔

    More partitions: more parallelism and headroom; fewer partitions: cheaper operations, better batching and faster recovery.

Done when you can

  • I can size partitions from measured throughput and consumer parallelism with headroom.

  • I know why increasing partitions breaks key ordering.