skip to content

What practical limits govern how many partitions a single broker (and the whole cluster) can host, and what breaks when you exceed them?

level: seniorimportance: must knowfreq 65%

answer

  1. replicas, not topics, are the cost
  2. ~3 fds per open segment → raise nofile
  3. ZK ~200k cap; KRaft → millions
  4. more partitions = slower failover + recovery
  5. size by consumer parallelism

basics

~10 s

Each partition replica is open files plus memory and adds controller/recovery work. A rough guide is up to ~1000-4000 partition-replicas per broker. Too many slows leader election, lengthens unclean-shutdown recovery, and increases end-to-end latency.

solid answer

~40 s

Partition count is bounded per broker and per cluster. Per broker, each partition *replica* consumes file descriptors (segment data + .index + .timeindex files), page-cache and heap, plus a fetch/replication thread share — historically people cap around 1000-4000 replicas per broker. Per cluster the old ZooKeeper limit was ~200k partitions; KRaft raises this substantially (millions of partitions targeted) because partition metadata lives in a replicated metadata log instead of ZooKeeper znodes. Exceeding limits causes: longer controlled and *unclean* shutdown recovery (each partition's log must be checked/reloaded), slower leader failover (more leaders to move when a broker dies), higher producer latency (more requests, larger metadata), and file-descriptor exhaustion. Mitigations: raise the `nofile` ulimit, right-size partitions to actual parallelism need, prefer KRaft, and watch `UnderReplicatedPartitions` and controller metrics.

go deeper

for a junior

Know partitions are the parallelism unit and that 'too many' hurts; you can add but not remove partitions.

for a middle

Account for fds, page cache, and the per-broker replica ceiling; size partitions from throughput/parallelism needs.

for a senior

Reason about failover and recovery time scaling with partition count, ZK vs KRaft limits, and monitoring signals.

for a principal

Set cluster-wide partition budgets, drive KRaft migration, and design topologies that survive broker loss within limits.

## Why partition count is a capacity dimension A Kafka topic is split into **partitions** — the unit of parallelism and ordering. Each partition is **replicated** RF times; every copy is a *replica*. The broker-level cost is driven by total **replicas hosted**, not topics. ## Per-broker costs of a replica 1. **File descriptors**: every open segment contributes a data file + `.index` + `.timeindex` = ~3 fds, and the active segment is always open. Thousands of partitions x multiple segments can blow past the default `nofile` ulimit (often 1024). You must raise it (tens/hundreds of thousands). 2. **Memory / page cache**: Kafka relies on the OS page cache for hot reads/writes. More partitions spread cache thinner; heap also holds per-partition state. 3. **Threads & requests**: replication fetchers, request handlers, and produce/fetch batching cost scales with active partitions; many tiny partitions mean many small I/Os instead of few large sequential ones. ## Cluster-level costs - **Controller / metadata**: in the legacy **ZooKeeper** mode, partition metadata sat in znodes and the controller did per-partition work on failover; clusters were practically limited to ~200,000 partitions and leader election after a controller failure could take many seconds. - **KRaft** (KIP-500) replaces ZooKeeper with an internal Raft **metadata log** and a quorum of controllers; metadata changes are incremental and replicated, raising the supported partition count by orders of magnitude (Kafka targets millions of partitions) and making failover near-constant time. ## What breaks when you over-partition - **Unclean/hard shutdown recovery**: on restart after a crash, the broker must scan/recover each partition log (checkpoint + index rebuild). Recovery time scales with partition count — thousands of partitions can mean minutes of downtime. `num.recovery.threads.per.data.dir` parallelizes this. - **Leader failover latency**: when a broker dies, all partitions it led must elect new leaders. More partitions = longer unavailability window and a metadata-update storm. - **End-to-end latency & throughput cliff**: too many partitions fragment sequential I/O and inflate request counts; producer/consumer metadata grows. - **fd exhaustion**: 'Too many open files' errors crash the broker. ## Rules of thumb (not hard limits) - A common operational ceiling is ~**1,000-4,000 partition-replicas per broker** depending on hardware and Kafka version. - An old Confluent guideline: keep **partitions <= 100 x brokers x RF** for the cluster. - Size partitions by required **consumer parallelism** and per-partition throughput target (e.g. each partition handles X MB/s; a topic at N MB/s needs ~N/X partitions), then verify the per-broker total stays under your ceiling. ## Monitoring Watch `UnderReplicatedPartitions`, `OfflinePartitionsCount`, controller `ActiveControllerCount`, fd usage, and recovery/log-load times. Plan with headroom for one or more broker failures (the survivors absorb the dead broker's leaders).

  • Why does KRaft dramatically raise the partition ceiling compared with ZooKeeper?
    ZooKeeper stored partition metadata as znodes and the controller pushed full per-partition updates on failover, bottlenecking at ~200k partitions. KRaft keeps metadata in a replicated Raft log, applies incremental deltas, and lets controllers fail over near-instantly, supporting orders of magnitude more partitions.
  • You suddenly need more partitions on an existing topic for parallelism. What's the catch?
    You can only increase partitions, never decrease. Adding partitions changes key→partition mapping (default hash partitioner uses partition count), so existing keys may move and per-key ordering across the change is not preserved. Plan partition count up front for keyed topics.

saying these in an interview costs you the question

  • Claiming there's a hard universal partition limit — it depends on hardware, version (ZK vs KRaft), and segment count.
  • Thinking topic count, not replica count, drives broker load.
  • Ignoring file-descriptor (nofile) limits.
  • Believing you can reduce a topic's partition count later.

context