skip to content

Cluster Sizing and Capacity Planning

Sizing brokers, partitions and disk from throughput and retention numbers, with headroom for page cache, network and file descriptors. A classic back-of-the-envelope interview exercise.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

How do you estimate the raw disk storage a Kafka topic (or cluster) will consume given its throughput, retention, and replication settings?

level: juniorimportance: must knowfreq 75%

answer

  1. throughput x retention x RF
  2. retention.ms vs retention.bytes (first wins)
  3. RF multiplies copies
  4. keep disks <60-70% full
  5. compaction = keys, not time

basics

~20 s

Storage = ingest rate (bytes/sec) x retention seconds x replication factor. So a 10 MB/s topic kept for 7 days at RF=3 needs roughly 10MB x 604800s x 3 = about 18 TB of disk across the cluster.

solid answer

~40 s

The core capacity formula is: disk = write_throughput (bytes/s) x retention_seconds x replication_factor. retention_seconds comes from retention.ms (time-based) or you cap it with retention.bytes (size-based). Replication multiplies because every message is stored once per replica (RF copies total). For a topic doing 10 MB/s, retained 7 days (604800 s), RF=3: 10e6 x 604800 x 3 = 18.1 TB cluster-wide, ~6 TB per copy. On top of raw data you add headroom: index/timeindex files, the active segment that isn't deleted until it rolls (segment.bytes / segment.ms), compression ratio (post-compression bytes on disk if producers compress), and a free-space buffer (keep disks <60-70% full so log cleaning and rebalances have room). Compacted topics size differently — they're bounded by key cardinality, not time.

go deeper

for a junior

Recall the formula: bytes/sec x retention seconds x replication factor. Know RF multiplies storage.

for a middle

Apply retention.ms vs retention.bytes, account for compression and headroom, compute per-broker disk.

for a senior

Reason about active-segment lag, free-space buffers for rebalance/recovery, and compacted-topic sizing.

for a principal

Define org-wide capacity standards, model growth/peak vs average, and codify utilization targets and headroom policy.

## The problem Kafka stores every message on disk as an immutable append-only log. Before provisioning a cluster you must predict how much disk you need so you don't run out (a full disk takes a broker offline and can cascade). ## Key terms - **Throughput / ingest rate**: how many bytes per second producers write to a topic. Measure the *post-batching, post-compression* bytes that actually land on disk if producers compress. - **Retention**: how long Kafka keeps data before deleting it. Controlled by `retention.ms` (time, default 7 days = 604800000 ms) and/or `retention.bytes` (max bytes *per partition*). Whichever limit is hit first triggers deletion of old log segments. - **Replication factor (RF)**: number of copies of each partition across brokers (`replication.factor`). RF=3 means three full copies exist, so disk usage is multiplied by 3. - **Segment**: the log is split into files of `segment.bytes` (default 1 GiB) or rolled after `segment.ms`. Retention deletes whole *closed* segments; the *active* segment is never deleted, so a low-traffic partition can hold data past its retention until the segment rolls. ## The formula ``` total_disk = avg_write_bytes_per_sec x retention_seconds x replication_factor ``` Worked example: 10 MB/s, retention 7 days, RF=3: - per-copy: 10e6 B/s x 604800 s = 6.048e12 B ≈ 6.05 TB - cluster total: x3 = 18.14 TB Per broker, divide by the number of brokers (assuming even partition spread): 18.14 TB / 6 brokers ≈ 3 TB each. ## What the formula omits (add headroom) 1. **Index files**: each segment has `.index` (offset→position) and `.timeindex`. Small (a few %) but real. 2. **Compression**: if producers set `compression.type` (e.g. lz4, zstd), disk holds the compressed bytes. Measure actual bytes-on-wire, not uncompressed app data. 3. **Active segment lag**: data lingers past retention until the active segment rolls — significant for many low-volume partitions. 4. **Free space buffer**: keep disks below ~60-70%. Replica reassignment, log compaction, and recovery all need scratch space; a 100%-full disk halts the broker. 5. **Replica catch-up**: a reassignment temporarily creates a 4th copy of moving partitions. ## Compacted topics are different For `cleanup.policy=compact`, retention isn't time x throughput — Kafka keeps the *latest value per key*, so size ≈ (number of distinct keys) x (avg record size) x RF, plus a 'dirty' buffer of un-compacted recent writes. ## Putting it together Provision: `(throughput x retention x RF) / target_utilization`, e.g. divide the 18.14 TB by 0.65 ≈ 28 TB of raw disk to keep utilization safe.

  • Why divide the result by a target utilization like 0.65 instead of provisioning exactly the computed size?
    Operations like replica reassignment, recovery, and log compaction need scratch space; a full disk takes the broker offline. Headroom (~30-40%) absorbs bursts, the active segment lag, and rebalance traffic.
  • How would the math change for a log-compacted topic?
    Compaction keeps the latest record per key, so size is bounded by distinct-key cardinality x avg record size x RF, not retention.ms x throughput. Time-based retention math doesn't apply to the compacted portion.

saying these in an interview costs you the question

  • Forgetting to multiply by replication factor (under-provisions by RF-times).
  • Using uncompressed application bytes when producers actually compress.
  • Sizing to 100% disk utilization with no headroom for rebalances/recovery.
  • Assuming retention.bytes is cluster-wide — it's per partition.

context

open as a page

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%

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.

open as a page

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

level: middleimportance: should knowfreq 55%

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.

open as a page

Why is OS page cache central to Kafka sizing, and how do you reserve page-cache and network headroom when provisioning brokers?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Kafka serves recent reads from the OS page cache (RAM) instead of disk, so consumers that keep up stay memory-fast. Give the JVM a modest heap (~5-6 GB) and leave most RAM free for page cache. For network, size NICs above peak replication + produce + fetch traffic, which RF multiplies.

open as a page

How do you decide the number of brokers in a cluster, accounting for replication factor, broker-failure headroom, and rebalance capacity?

level: principalimportance: should knowfreq 40%

basics

~20 s

Pick brokers so total disk/network/partition load fits with each broker under ~60-70% utilization, you have at least RF brokers (plus spare for min.insync.replicas), and the cluster keeps working when N brokers fail. Survivors must absorb the dead brokers' load and re-replication.

open as a page