skip to content

Walk me through the per-phase timing metrics in a Kafka broker's request lifecycle. Which ones do you read first to localize where p99 latency is being spent?

level: middleimportance: must knowfreq 78%

answer

  1. RequestMetrics phases are additive to TotalTimeMs
  2. Queue / Local / Remote / Throttle / ResponseSend
  3. RemoteTime = waiting on others (acks=all, purgatory)
  4. Read 99th/999th, not Mean
  5. split first, then drill

basics

~20 s

Kafka breaks each request's total time into phases: queue, local processing, remote (waiting on other brokers), throttle, and response send. You read these phase metrics (RequestQueueTimeMs, LocalTimeMs, RemoteTimeMs, ResponseQueueTimeMs, ResponseSendTimeMs) to see which phase dominates.

solid answer

~40 s

Every broker request emits a per-request-type histogram under kafka.network:type=RequestMetrics broken into additive phases: RequestQueueTimeMs (waiting in the request queue for an I/O thread), LocalTimeMs (leader-local processing), RemoteTimeMs (time blocked waiting on other brokers — e.g. acks=all replication or a fetch purgatory delay), ThrottleTimeMs (quota throttling), ResponseQueueTimeMs, and ResponseSendTimeMs, summing to TotalTimeMs. To localize a p99 spike I pull the 99th percentile of each phase for the offending request type (Produce / FetchConsumer / FetchFollower). High RequestQueueTime points at I/O-thread or queue saturation; high LocalTime at GC, lock contention, or disk; high RemoteTime at slow followers or purgatory waits; high ResponseSendTime at network/TCP backpressure. This phase split is the single most important first cut.

go deeper

for a junior

Know that Kafka splits each request's time into phases and that you read those phases to find where latency is spent.

for a middle

Name the phases, know they sum to TotalTimeMs, and map each high phase to a likely root cause.

for a senior

Use the phase split as the first diagnostic cut, distinguish expected RemoteTime (acks=all, long-poll) from pathological, and correlate per-broker percentiles.

for a principal

Design dashboards/alerts around per-phase p99/p999 per request type per broker, and teach teams to localize before drilling into GC/disk/network subsystems.

## The problem "p99 latency is high" is useless on its own — you must know *where in the request's journey* the time goes. Kafka instruments this for you. ## What a broker request is Clients (producers, consumers) and other brokers (followers fetching to replicate) all send **requests** to a broker over the network: ProduceRequest, FetchRequest, MetadataRequest, etc. The broker processes each one and sends a response. Kafka measures the wall-clock time each request spends in each internal stage. ## The thread pipeline (why the phases exist) 1. **Network threads** (`num.network.threads`) read bytes off the socket and place a parsed request onto a shared **request queue**. 2. **I/O / request-handler threads** (`num.io.threads`) pull from that queue and actually do the work (append to log, read from log). 3. Some requests cannot complete immediately and are parked in **purgatory** (a delayed-operation holding area) — e.g. a Produce with `acks=all` waits there until enough followers replicate; a Fetch with `fetch.min.bytes` waits until enough data accumulates or `fetch.max.wait.ms` elapses. 4. The response is queued, then a network thread writes it back to the socket. ## The additive phase metrics Under JMX `kafka.network:type=RequestMetrics,name=<Phase>,request=<Produce|FetchConsumer|FetchFollower|...>`, each phase is a histogram (Mean, 99thPercentile, 999thPercentile, Max): - **RequestQueueTimeMs** — time the request sat in the request queue before an I/O thread picked it up. High = not enough I/O threads, or I/O threads blocked. - **LocalTimeMs** — time the I/O thread spent doing leader-local work (validating, appending to / reading from the log, page-cache access). High = GC pauses, disk stalls, lock contention, compression/validation cost. - **RemoteTimeMs** — time the request spent **blocked waiting on something external**: replication acks from followers (Produce acks=all), or sitting in fetch purgatory waiting for data. This is *expected* to be non-zero for acks=all and long-poll fetches; it's a problem when it spikes beyond your config's intent. - **ThrottleTimeMs** — time added by quota throttling (client/user quotas). - **ResponseQueueTimeMs** — time the finished response waited for a network thread. - **ResponseSendTimeMs** — time spent actually writing the response to the socket. High = network saturation or a slow/backpressured client TCP connection. - **TotalTimeMs** = sum of all of the above. Because they sum to the total, **comparing the p99 of each phase tells you which stage owns the latency.** That is the first cut of every tail-latency investigation. ## Worked example Produce p99 = 250 ms. You pull phases: RequestQueueTime p99 = 5 ms, LocalTime p99 = 8 ms, **RemoteTime p99 = 230 ms**, send = 3 ms. Conclusion: the leader processes fast; the time is spent waiting for follower replication (acks=all) — investigate slow/lagging followers or ISR, not broker CPU. Counter-example: RemoteTime low but **RequestQueueTime p99 = 200 ms** → I/O threads can't keep up; check `num.io.threads`, the request-handler idle percent, and whether I/O threads are blocked on GC or disk. ## Edge cases - RemoteTime being high for FetchConsumer is often *normal* (long-poll waiting for `fetch.min.bytes`). Don't treat expected purgatory waits as a bug — compare against `fetch.max.wait.ms`. - Mean can look fine while p99/p999 are terrible; always read the high percentiles, not the mean. - Percentiles are per-broker; a cluster-wide p99 spike on one broker often means one bad node.

  • If RemoteTimeMs is high on Produce but low on Fetch, what does that suggest?
    Produce RemoteTime is dominated by acks=all replication waits, so high Produce RemoteTime points at slow followers / ISR replication lag rather than fetch purgatory. Look at replica fetch lag, follower GC, and ISR membership.
  • Why might RemoteTimeMs be legitimately high for FetchConsumer requests and not indicate a problem?
    Consumer fetches long-poll: they wait in purgatory up to fetch.max.wait.ms for fetch.min.bytes of data. On low-traffic topics that wait is expected and intentional, so high RemoteTime there is by design, not a stall.

saying these in an interview costs you the question

  • Treating TotalTimeMs as a single opaque number instead of splitting into phases
  • Assuming high RemoteTimeMs always means a broker problem (it's normal for acks=all and long-poll fetches)
  • Looking only at Mean and ignoring 99th/999th percentiles
  • Confusing RequestQueueTime (before processing) with ResponseQueueTime (after processing)

context