skip to content

What are the per-partition costs that make over-partitioning harmful, and what cluster limits should you weigh when choosing a high partition count?

level: seniorimportance: should knowfreq 55%

answer

  1. FDs: segment + .index + .timeindex
  2. x replicas x partitions
  3. ulimit -n -> Too many open files
  4. failover = elect a leader per partition
  5. KRaft raised the ceiling, not infinite

basics

~20 s

Each partition costs open file handles (log segments + index files), broker memory, replication threads, and controller/metadata load. More partitions also slow leader-election failover and lengthen rebalances. Open file descriptor limits and end-to-end latency cap how many partitions a cluster can hold.

solid answer

~50 s

A partition isn't free. Per partition the broker holds open file handles for each active log segment plus its offset and time indexes (so partitions × replicas × segments × ~3 files of open FDs), broker heap and page-cache pressure, a replica-fetcher workload, and an entry in cluster metadata that the controller must manage. Total partitions across a cluster are bounded by: OS open-file-descriptor limits (`ulimit -n`), broker memory, and — critically — failover time. When a broker dies, the controller must elect new leaders for every partition that broker led; with tens of thousands of partitions per broker this election and metadata propagation can take many seconds, raising unavailability windows. More partitions also lengthen consumer-group rebalances and increase end-to-end latency for replication. KRaft raised practical ceilings versus ZooKeeper, but the guidance remains: size for needed parallelism plus modest headroom, not 'as many as possible'.

go deeper

for a junior

Know that more partitions cost more files, memory, and slower failover — they aren't free.

for a middle

List the concrete per-partition costs (FDs, replication threads, metadata) and the ulimit failure mode.

for a senior

Trade parallelism benefits against failover/rebalance time and total-partition cluster limits when sizing.

for a principal

Set org-wide partition budgets, account for KRaft vs ZooKeeper ceilings, and design topic taxonomies that keep total partition counts sane.

## Why a partition costs resources A **partition** is a physical append-only log on disk, replicated across brokers. Each partition replica incurs: - **Open file descriptors.** A partition's log is split into **segments**; each active segment is an open file, plus its `.index` (offset index) and `.timeindex` (time index) files. So a single partition replica holds roughly 3 open files for its active segment, more if older segments are still open. Multiply by partitions × replication factor across the broker. Hitting the OS `ulimit -n` (`Too many open files`) crashes brokers — a classic over-partitioning failure. - **Memory.** Each partition consumes broker heap for in-memory index structures and producer/fetch state, and competes for OS page cache (which is what makes Kafka fast). Thousands of partitions dilute page-cache effectiveness. - **Replication work.** Each follower replica runs through replica-fetcher threads; more partitions = more fetch requests, more I/O, more network. - **Metadata / controller load.** The controller tracks leadership and ISR (in-sync replicas) for every partition. More partitions = a larger metadata footprint to store, replicate, and propagate. ## Cluster-level limits to weigh 1. **Open file descriptors** — set `ulimit -n` high (tens/hundreds of thousands), but it's still a ceiling. Partitions × replicas × open segments must stay under it. 2. **Failover / leader-election time.** When a broker fails, the controller elects new leaders for *every partition that broker was leading*. Under ZooKeeper this was roughly linear and could take many seconds-to-minutes at very high partition counts (historically a few thousand per broker was a practical concern); the partitions are unavailable until election completes. KRaft (the ZooKeeper-free metadata mode) improved this substantially, pushing supported per-cluster partition counts into the millions, but failover is still not instant. 3. **Rebalance time.** Consumer-group rebalances scale with partition count; huge topics make every rebalance (and every deploy/scaling event) slower. 4. **End-to-end latency.** More partitions to replicate can raise tail latency for produce acks (`acks=all`). ## Right-sizing Balance against the *benefits* of partitions (parallelism ceiling, throughput). Practical rule of thumb: estimate required throughput and consumer parallelism, pick partitions to cover that plus headroom (since you can't shrink and growth breaks key routing), but resist 'just make it 1000 to be safe' — that tax is paid continuously by every broker and every rebalance. ## Edge cases - Idle/low-traffic topics still cost FDs and metadata even with no data flowing. - Many *small* topics aggregate the same way as few *large* topics — the cluster cares about *total partitions*, not per-topic counts. - Compacted topics keep more segments around, raising open-FD pressure further.

  • A broker is crashing with 'Too many open files'. How is partition count implicated?
    Each partition replica keeps open file handles for its active log segment plus the offset (.index) and time (.timeindex) index files. With many partitions × replicas the total open FDs can exceed the OS ulimit -n, crashing the broker. Fixes: raise ulimit -n, reduce total partitions, or rebalance partitions off the broker.
  • Why does a very high partition count slow recovery when a broker dies?
    The controller must elect a new leader for every partition the failed broker was leading and propagate the new metadata. That work scales with partition count, so more partitions means a longer window where those partitions are leaderless and unavailable. KRaft reduced but did not eliminate this cost.

saying these in an interview costs you the question

  • Claiming partitions are essentially free so you should always over-provision heavily.
  • Ignoring the open-file-descriptor cost (segment + index files per partition replica).
  • Saying failover time is independent of partition count.
  • Forgetting that the cluster cares about TOTAL partitions across all topics, not per-topic.

context