02 / 05

Average event latency is healthy but p99 latency spikes under load. What evidence would you collect before optimizing?

Difficulty: 8/10
Performance incident

Measure end-to-end latency per stage with percentile histograms, correlate spikes with broker, client and infrastructure saturation signals, and only then form a hypothesis

Averages hide exactly this problem, so the first step is to make the measurement trustworthy. I would define what latency means (producer send to consumer handler start, or to downstream effect), record it as a histogram so p50, p95, p99 and p99.9 are available, and break it into stages: time in the producer before the request leaves (record-queue-time), request round trip including replication (request-latency), time from broker append to consumer receipt, and consumer queueing before processing. I carry a producedAt timestamp in a header and compare it with the broker's log-append time and the consumer clock, being careful about clock skew between hosts. I also check that the load generator is not hiding the problem through coordinated omission, where a stalled sender stops generating the slow samples.

Then I line up the spike windows against saturation signals using a USE-style view (utilization, saturation, errors) on every layer. Broker: request handler and network processor idle percentages, request queue time versus local, remote and response-send time on produce and fetch requests, disk await, page cache pressure, GC pauses, ISR shrink and expand events, and quota throttle times. Client: batch size and linger behavior, buffer exhaustion, retries and rebalances. Consumer: lag over time, poll-idle ratio, processing time, and downstream call latency. Infrastructure: CPU throttling, network drops, noisy neighbors, cross-AZ hops and TLS cost. The goal is a short list of hypotheses ranked by evidence, such as acks=all stalling on one slow follower, lagging consumers forcing disk reads and evicting hot data from page cache, leader skew, or periodic GC, each with an experiment that would confirm or refute it. I would change nothing until one hypothesis is supported, and I would change one thing at a time.

javascript
  1. 1

    Trade-off: a full histogram per stage costs instrumentation effort and metric cardinality, but without it you only know that p99 is bad, not where. I accept the cost on the critical paths and sample elsewhere.

  2. 2

    Trade-off: tracing every message gives exact causality but adds overhead and volume. I use trace sampling plus always-on histograms, and turn on detailed tracing only around spikes.

  3. 3

    Common mistake: optimizing from averages or from p50 dashboards. A system with a 5 ms mean and a 900 ms p99 can have a perfectly healthy average.

  4. 4

    Common mistake: tuning producer linger.ms or batch.size first. If the stall is replication, disk or GC on one broker, batching changes will not fix it and can add latency.

  5. 5

    Common mistake: ignoring that lagging consumers reading old data push reads to disk, which can evict recent data from page cache and slow everyone's fetches, producing a latency spike that looks unrelated to the consumer.

  6. 6

    Common mistake: trusting cross-host timestamps blindly. Producer, broker and consumer clocks differ, so use NTP-synced hosts and compare deltas on the same host where possible.

  7. 7

    Common mistake: changing several settings at once, which makes it impossible to know what helped. Change one variable, rerun the same load, compare histograms.

  8. 8

    Verification: reproduce under a controlled load test with failure injection (kill a broker, throttle a disk) so the fix is validated against the real spike pattern rather than assumed.

  9. 9

    Version note: metric names above exist in current clients and brokers, but broker request metric breakdowns and consumer metric names differ in the new consumer implementation introduced with the KIP-848 protocol in Kafka 4.0. Verify names on your version.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.