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?
answer
- RequestMetrics phases are additive to TotalTimeMs
- Queue / Local / Remote / Throttle / ResponseSend
- RemoteTime = waiting on others (acks=all, purgatory)
- Read 99th/999th, not Mean
- split first, then drill
basics
~20 sKafka 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 sEvery 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
Know that Kafka splits each request's time into phases and that you read those phases to find where latency is spent.
Name the phases, know they sum to TotalTimeMs, and map each high phase to a likely root cause.
Use the phase split as the first diagnostic cut, distinguish expected RemoteTime (acks=all, long-poll) from pathological, and correlate per-broker percentiles.
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)