skip to content

Broker Threads, Page Cache and Disk Tuning

Broker-side tuning: thread pools, request queues, page cache versus flush settings, filesystem and disk layout, and heap sizing. Interviewers ask because Kafka's tuning advice inverts normal JVM instincts.

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

questions

5

What do num.network.threads and num.io.threads control on a Kafka broker, and how do you decide how many to set?

level: juniorimportance: must knowfreq 70%

answer

  1. network=read/write sockets, io=do the work
  2. defaults 3 and 8
  3. Idle% JMX gauges drive sizing
  4. queued.max.requests sits between them
  5. io near cores/disks, network smaller

basics

~10 s

num.network.threads handle reading requests off and writing responses onto network sockets; num.io.threads do the actual work (reading/writing the log on disk). Default 3 and 8. Raise them if request-handler/network idle time drops.

solid answer

~40 s

A Kafka broker uses two thread pools. Network threads (num.network.threads, default 3) run a Reactor loop: they accept socket connections, read raw request bytes, and write response bytes back, then hand the parsed request to a shared request queue. I/O threads, also called request handlers (num.io.threads, default 8), pull from that queue and do the heavy work: appending to the partition log, reading data for fetches, handling metadata. You size them from utilization metrics, not guesswork: the JMX gauges kafka.network:type=SocketServer,name=NetworkProcessorAvgIdlePercent and kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent. If idle percent trends toward 0, that pool is saturated and you scale it up. I/O threads are typically set near the number of disks/cores; network threads stay smaller. Over-provisioning wastes context-switching and memory.

go deeper

for a junior

Know the two pools, their defaults (3 / 8), and that network=sockets, io=disk work.

for a middle

Name the idle-percent JMX gauges and use them to decide which pool to scale.

for a senior

Reason about queued.max.requests backpressure and dynamic (no-restart) reconfiguration from live metrics.

for a principal

Frame thread sizing within end-to-end broker capacity: TLS cost on network threads, disk fsync stalling I/O threads, and cores/disks as the real ceiling.

## The two thread pools A Kafka broker is a network server that persists data to disk. To keep the network layer and the disk/CPU layer from blocking each other, it splits work across two thread pools connected by a bounded queue. **Network threads — `num.network.threads` (default 3).** Internally these are the `Processor` threads of the `SocketServer`. They run an NIO event loop (a Reactor pattern): one `Acceptor` accepts new TCP connections and round-robins them to processors; each processor reads inbound request bytes off its sockets, assembles complete requests, and pushes them onto a shared **request queue**. After an I/O thread produces a response, the processor writes the response bytes back to the socket. Crucially, network threads do **no** log or disk work — they must never block, or all connections they own stall. **I/O threads — `num.io.threads` (default 8).** These are the `KafkaRequestHandler` threads (the "request handler pool"). Each one loops: take a request from the request queue, dispatch it to `KafkaApis`, perform the real work — appending records to the active log segment, reading bytes for a fetch, serving metadata/offset/group requests — then enqueue the response for the originating network thread. **The queue between them — `queued.max.requests` (default 500).** Bounds how many parsed-but-unhandled requests can pile up. When full, network threads stop reading new requests (backpressure), which protects broker heap from unbounded buffering. ## How to size them — measure, don't guess Kafka exposes idle-percent gauges via JMX: - Network: `kafka.network:type=SocketServer,name=NetworkProcessorAvgIdlePercent` (0.0–1.0). - I/O: `kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent`. Interpretation: idle near 1.0 = pool mostly waiting (fine/over-provisioned); idle trending toward 0 = pool saturated and becoming the bottleneck → add threads. A common operational rule of thumb is to keep both above ~0.3. Increase `num.io.threads` when the request-handler idle drops (often the first bottleneck on disk-heavy workloads); increase `num.network.threads` when the network-processor idle drops (high connection counts / TLS / many small requests). ## Sizing heuristics - I/O threads: scale toward the number of CPU cores and number of data directories (disks); values like 8–16 are common on multi-disk brokers. - Network threads: usually smaller than I/O threads; TLS termination and very high connection counts push it up. - Both are dynamically updatable cluster-wide via `kafka-configs.sh --alter --entity-type brokers` without a restart, which lets you tune from metrics in production. ## Edge cases / pitfalls - Setting them absurdly high doesn't help: extra threads add context-switch overhead and memory; the disk or NIC becomes the real ceiling. - A blocked I/O thread (e.g., slow disk fsync) drains the request queue slowly, request-handler idle drops, and producer/consumer latency spikes — the metric tells you which pool to fix. - Network and I/O idle are independent: a TLS-heavy cluster can saturate network threads while I/O threads sit idle.

  • Which exact metrics tell you a pool is saturated?
    NetworkProcessorAvgIdlePercent and RequestHandlerAvgIdlePercent — as idle trends to 0 the pool is the bottleneck; keep them comfortably above ~0.3.
  • Can you change these without restarting the broker?
    Yes. Both are dynamic broker configs; alter them cluster-wide via kafka-configs.sh --entity-type brokers and they apply live.

saying these in an interview costs you the question

  • Saying network threads write to disk — they only do socket I/O.
  • Claiming 'more threads = more throughput' regardless of disk/NIC limits.
  • Confusing num.io.threads with producer/consumer client threads.
  • Not knowing any metric to justify the value.

context

open as a page

Why does Kafka rely on the OS page cache instead of flushing every write to disk, and what is the role of log.flush.interval.messages / log.flush.interval.ms?

level: middleimportance: must knowfreq 65%

basics

~20 s

Kafka writes to the OS page cache and lets the OS flush to disk in the background, which is fast. Durability comes from replication, not fsync. log.flush.interval.* would force periodic fsyncs, but Kafka leaves them effectively off by default.

open as a page

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%

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.

open as a page

Compare JBOD vs RAID for Kafka brokers and explain why XFS is commonly recommended over ext4.

level: seniorimportance: should knowfreq 40%

basics

~20 s

JBOD gives Kafka each disk separately via log.dirs, so failures are isolated and replication handles redundancy; RAID (often RAID10) adds redundancy and even I/O but costs capacity/write speed. XFS is favored over ext4 for better performance under Kafka's large sequential-write workload.

open as a page

Explain zero-copy / sendfile in Kafka: what it optimizes, when it applies, and what disables it.

level: seniorimportance: should knowfreq 45%

basics

~20 s

Zero-copy uses the OS sendfile() call to stream bytes from the page cache straight to the network socket, skipping copies into and out of the JVM heap. It makes consumer fetches cheap — but only for plaintext data Kafka doesn't have to transform.

open as a page