skip to content

A consumer reports high TotalTimeMs and you see RequestQueueTimeMs dominating with a large RequestQueueSize. What does this indicate and how do you remediate it?

level: seniorimportance: should knowfreq 55%

answer

  1. RequestQueueTime dominant → handler pool starved
  2. num.io.threads (handlers), num.network.threads
  3. RequestHandlerAvgIdlePercent ≈ 0 confirms
  4. queued.max.requests = queue capacity, not capacity
  5. fix slow downstream phase before adding threads

basics

~10 s

Requests are piling up in the request queue because there aren't enough I/O (request handler) threads to process them. You typically increase num.io.threads, and check for slow downstream phases that back up the handlers.

solid answer

~40 s

A dominant RequestQueueTimeMs plus a large RequestQueueSize means requests are accepted off the socket faster than the request-handler (I/O) thread pool can process them — the queue is backing up. The pool is sized by num.io.threads (handlers, default 8) and the queue capacity by queued.max.requests (default 500). First confirm via the RequestHandlerAvgIdlePercent metric: if it's near 0, handlers are saturated. Remediation: increase num.io.threads (rule of thumb a few × cores), and increase num.network.threads if network threads are also saturated (NetworkProcessorAvgIdlePercent low). But also investigate *why* handlers are busy: if LocalTimeMs is high (slow disk) or RemoteTimeMs is high, adding threads only masks a downstream bottleneck. So fix the slow phase first, then scale threads to match load.

go deeper

for a junior

Know that requests wait in a queue and too few worker threads cause queue buildup.

for a middle

Distinguish network threads from I/O/request-handler threads and name num.io.threads / num.network.threads.

for a senior

Confirm via RequestHandlerAvgIdlePercent, and fix the underlying slow phase before scaling threads.

for a principal

Reason about disk contention vs thread count tradeoffs and set capacity policy / dynamic reconfiguration across the fleet.

## The threading model A Kafka broker has two thread pools on the request path: - **Network threads** (`num.network.threads`, default 3) — the `Processor`s that read bytes off client sockets and write responses. They do *not* do the business logic; they hand requests to a shared **request queue**. - **I/O / request-handler threads** (`num.io.threads`, default 8) — `KafkaRequestHandler`s that pull from the request queue and actually process the request (append to log, read from log, etc.). The request queue between them has a bounded capacity set by `queued.max.requests` (default 500). When it fills, network threads stop reading new requests (backpressure to clients). ## What the symptom means `RequestQueueTimeMs` measures how long a request waited in that queue before a handler picked it up. If it dominates TotalTimeMs and `RequestQueueSize` (a gauge) is large, requests are arriving faster than handlers drain them — the handler pool is the bottleneck. ## Confirming the diagnosis Look at **RequestHandlerAvgIdlePercent** (`kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent`). It ranges 0–1; near 0 means handlers are 100% busy. Similarly **NetworkProcessorAvgIdlePercent** for network threads. If handler idle is near 0, the handler pool is starved. ## Remediation, in order 1. **Find out why handlers are slow.** Decompose the *rest* of TotalTimeMs. If LocalTimeMs is high, handlers are blocked on slow disk/page-cache — adding more threads just creates more disk contention. If RemoteTimeMs is high, the bottleneck is replication or purgatory, not threads. Fix the real cause first. 2. **If handlers are genuinely CPU/throughput-bound**, raise `num.io.threads` (commonly 2–4× CPU cores). These are dynamically reconfigurable in modern Kafka without a restart. 3. **If network threads are saturated** (NetworkProcessorAvgIdlePercent low, ResponseSendTime/ResponseQueueTime high), raise `num.network.threads`. 4. **Don't blindly raise `queued.max.requests`** — a bigger queue just hides latency by letting more requests wait longer; it does not add processing capacity. ## Edge cases - A request storm of expensive requests (e.g. large metadata or many small produces) can saturate handlers even with healthy disk. - On disk-bound brokers, more I/O threads can *worsen* tail latency by increasing concurrent disk contention.

  • Why might increasing num.io.threads make tail latency worse on some brokers?
    If handlers are blocked on disk I/O (high LocalTimeMs), more threads mean more concurrent disk access and page-cache contention, which can increase p99 rather than reduce it. Threads help only when handlers are CPU/throughput-bound, not when they're waiting on a downstream resource.
  • What does RequestHandlerAvgIdlePercent near 1.0 tell you about a queue-time problem?
    Handlers are mostly idle, so the request queue is NOT the bottleneck. A high RequestQueueTimeMs with idle handlers would point elsewhere (e.g. brief bursts), and you'd look at network threads or arrival patterns instead of growing the I/O pool.

saying these in an interview costs you the question

  • Recommending raising queued.max.requests to fix latency (it only lets requests wait longer).
  • Always adding I/O threads without checking whether LocalTimeMs/RemoteTimeMs is the real bottleneck.
  • Confusing network threads (socket I/O) with I/O/request-handler threads (request processing).

context