What do num.network.threads and num.io.threads control on a Kafka broker, and how do you decide how many to set?
answer
- network=read/write sockets, io=do the work
- defaults 3 and 8
- Idle% JMX gauges drive sizing
- queued.max.requests sits between them
- io near cores/disks, network smaller
basics
~10 snum.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 sA 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
Know the two pools, their defaults (3 / 8), and that network=sockets, io=disk work.
Name the idle-percent JMX gauges and use them to decide which pool to scale.
Reason about queued.max.requests backpressure and dynamic (no-restart) reconfiguration from live metrics.
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.