02 / 05

Why is creating thousands of partitions not free?

Difficulty: 4/10
Large-scale design, Partitioning, Capacity planning

The Hidden Costs of Thousands of Partitions

Creating thousands of partitions is not free because each partition carries metadata, file handles, memory, and replication overhead. Every partition has at least one log segment with a .log file, an .offset index file, and a .time index file on disk. Each partition also has replication state: the leader tracks the ISR, the followers fetch from the leader, and the controller manages leadership and reassignment. The metadata for each partition is stored in the controller's metadata (ZooKeeper or KRaft), and the cluster's metadata grows with the number of partitions. The controller must process leader elections, ISR changes, and partition reassignments for every partition, so more partitions mean more controller work, especially during broker failures when many partitions need new leaders. Memory usage also grows: the broker caches partition metadata, and the consumer caches partition state. File handles are a limited resource; each open segment consumes a file handle, and a broker with thousands of partitions can exhaust the file handle limit. The trade-off is between parallelism and overhead: more partitions give more parallelism but cost more in metadata, memory, file handles, controller work, and recovery time.

The mechanism of the overhead has several parts. First, the controller: when a broker fails, the controller must elect new leaders for all partitions that had their leader on that broker. With thousands of partitions, this can take seconds to minutes, during which the cluster is partially unavailable for those partitions. Second, replication: each partition's leader must replicate to its followers. More partitions mean more replication streams, more network connections, and more disk I/O. Third, the log segment files: Kafka creates a new segment when the current one reaches a size or time limit. Each segment has index files, and the broker must keep track of them. Fourth, the consumer: each consumer in a group must track its position for each assigned partition, and the group coordinator must manage the assignment. With thousands of partitions, the group coordinator and the consumers have more state to manage. Fifth, the metadata: the metadata log (in KRaft) or ZooKeeper (in older versions) stores the configuration for each partition. The metadata grows linearly with the number of partitions, and the controller's ability to process metadata changes is a bottleneck. Version note: KRaft was specifically designed to scale to millions of partitions by removing ZooKeeper's limits, but the per-broker overhead of partitions remains. Kafka 3.x has improved partition scalability, but the practical limit per broker is still a few thousand partitions for most workloads.

A common mistake is to create a huge number of partitions for a topic that does not need it, such as a topic with a single consumer group that cannot process more than a few partitions' worth of data. Another mistake is to ignore the recovery time: a broker with 4,000 partitions takes longer to restart and rejoin the cluster than a broker with 400 partitions. A third mistake is to ignore the file handle limit: the default ulimit on many Linux systems is 1,024 or 4,096, which is easily exhausted by a broker with thousands of partitions. The trade-off is between the flexibility to scale and the operational cost. A good rule of thumb is to keep the number of partitions per broker below 4,000, and to monitor the controller's processing time, the broker's memory and file handle usage, and the recovery time after a broker failure. If you need more partitions, consider adding brokers or using a different topic design. Version note: Kafka 3.x with KRaft has better metadata scalability, but the partition count still affects recovery time, replication overhead, and consumer assignment complexity. Always test with a realistic partition count and workload.

javascript
  1. 1

    Each partition has log segments, index files, and replication state on disk.

  2. 2

    Controller work grows with partition count: leader elections, ISR changes, reassignments.

  3. 3

    Memory usage grows: brokers and consumers cache partition metadata.

  4. 4

    File handles are limited; thousands of partitions can exhaust them.

  5. 5

    Recovery time after a broker failure grows with partition count.

  6. 6

    Metadata grows linearly with partitions; the controller is a bottleneck.

  7. 7

    Keep partitions per broker below ~4,000 for most workloads.

  8. 8

    KRaft improves metadata scalability but per-broker overhead remains.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.