skip to content

How should you size the JVM heap and choose a garbage collector for a Kafka broker, and why is a relatively small heap usually correct?

level: seniorimportance: should knowfreq 50%

answer

  1. Data off-heap in files → small heap (~5-6 GB)
  2. Rest of RAM = page cache
  3. -Xms == -Xmx, no resize
  4. G1GC default, MaxGCPauseMillis~20
  5. Long GC pause → drop from ISR

basics

~20 s

Give the broker a modest heap (commonly ~5-6 GB) and leave most RAM to the OS page cache, since Kafka stores data off-heap in files. Use the G1 garbage collector (the default) to keep pauses low.

solid answer

~50 s

A Kafka broker's data lives in log files cached by the OS page cache, not on the JVM heap, so the heap only needs to hold request buffers, metadata, replication state and connection objects. The standard guidance is a fairly small fixed heap — Confluent recommends around 5-6 GB — with -Xms equal to -Xmx to avoid resizing. The rest of system RAM should be left free for page cache, which is what actually accelerates Kafka; over-sizing the heap steals RAM from cache and lengthens GC pauses. The recommended collector is G1GC (the JVM default on modern Java and what Kafka ships with), which targets bounded pause times via -XX:MaxGCPauseMillis (default ~20ms in Kafka's tuning). You watch GC via the broker's GC logs and JMX, and alarm on long pauses, since a long stop-the-world pause can drop the broker out of the ISR or trigger session timeouts. Heap pressure usually means too many partitions/connections, not too little RAM.

go deeper

for a junior

Know the heap is small (~5-6 GB) because data lives in files/page cache, and Kafka uses G1GC.

for a middle

Explain -Xms=-Xmx and that spare RAM must go to page cache, not heap.

for a senior

Tie GC pause targets (MaxGCPauseMillis) to ISR membership and diagnose heap pressure as a partition/connection problem.

for a principal

Design broker capacity holistically: heap vs page-cache split, GC choice (G1/ZGC), swappiness, and partition-count limits across the fleet.

## Why the heap is small A common surprise: Kafka brokers run very large data volumes on a **small** JVM heap. The reason is architectural — the actual message data is stored in **log segment files** and cached by the **OS page cache** (kernel RAM), entirely **off** the JVM heap. The heap is only used for: - network request/response buffers and the request queue, - per-connection and per-partition metadata and replication/fetch state, - index structures, controller/metadata in KRaft, and assorted bookkeeping. None of that scales with total data stored; it scales with **partitions, connections, and request volume**. So a modest heap suffices regardless of how many TB the broker holds. ## Recommended sizing - **Heap size:** a fixed, modest value — Confluent's guidance is roughly **5–6 GB** for typical brokers; very large/partition-heavy brokers may go higher, but doubling the heap is rarely the right fix. - **Set `-Xms == -Xmx`** so the JVM doesn't resize the heap at runtime (resizing causes pauses and fragmentation). - **Leave the rest of RAM for page cache.** This is the crucial point: if a box has 64 GB, you give the broker ~6 GB heap and let ~50+ GB serve as page cache. Page-cache RAM is what keeps reads off disk and powers zero-copy; spending it on an oversized heap *hurts* performance and lengthens GC. ## Garbage collector - **G1GC (Garbage-First)** is the recommended and default collector for Kafka on modern JDKs (Kafka's startup scripts configure G1). G1 is a low-pause, region-based collector that tries to meet a **pause-time target** rather than minimizing total GC time. - Kafka's tuned defaults include **`-XX:MaxGCPauseMillis=20`** (aim for ~20 ms pauses) and `-XX:InitiatingHeapOccupancyPercent=35` (start concurrent marking earlier). These bias toward short, predictable pauses over raw throughput — correct for a latency-sensitive broker. - On the newest JDKs some operators evaluate **ZGC** for sub-millisecond pauses on large heaps, but G1 remains the standard, well-tested choice for typical broker heaps. ## Why pauses matter operationally A long **stop-the-world** GC pause freezes the broker. Consequences: - The broker can miss its session/heartbeat deadlines and be considered failed, or - fall out of the **ISR** (in-sync replica set) because it stops fetching/acking, triggering leadership churn and replication catch-up. So GC tuning targets *pause predictability*, not throughput. You enable GC logging (`-Xlog:gc*`) and monitor via JMX, alerting on pause spikes and high promotion/allocation rates. ## Diagnosing heap pressure If a broker shows GC pressure or OOMs, the usual root cause is **too many partitions or connections** (each costs heap), not insufficient RAM. The right responses are reducing partition count per broker, capping connections, or scaling out — not blindly enlarging `-Xmx`, which steals page-cache RAM. ## Edge cases / pitfalls - Giving Kafka a huge heap 'to be safe' is an anti-pattern: it starves page cache and produces longer GC pauses — strictly worse. - Swapping is deadly: set `vm.swappiness` low (e.g., 1) so the kernel doesn't swap out broker memory or page cache under pressure. - Heap need grows with partition count; brokers hosting tens of thousands of partitions need more heap and careful GC tuning.

  • Why not just give the broker a 32 GB heap to be safe?
    It steals RAM from the OS page cache (which is what actually speeds Kafka), and a bigger heap means longer GC pauses — both hurt. Heap should stay small; spare RAM goes to cache.
  • What operational failure can a long GC pause cause?
    The broker can fall out of the ISR or miss session/heartbeat deadlines, causing leadership churn and replication catch-up — hence tuning for bounded pauses (MaxGCPauseMillis).
  • If a broker has GC/heap pressure, what's the usual real cause?
    Too many partitions or connections per broker (each consumes heap), not too little RAM — fix by reducing partitions/connections or scaling out.

saying these in an interview costs you the question

  • Sizing the broker heap to hold the data — data is off-heap in page cache.
  • Recommending a very large heap 'for safety' — it hurts cache and GC.
  • Not knowing G1GC is the default/recommended collector.
  • Ignoring that long GC pauses can drop a broker from the ISR.
  • Leaving vm.swappiness high so the OS swaps the broker.

context