Monitoring and Performance Tuning
Measuring and tuning Kafka: broker and client JMX metrics, request-latency breakdowns, lag monitoring, and throughput tuning on both sides. Interviewers use it for 'the cluster is slow, what do you look at' questions.
part ofApache Kafkaoverview, primer and where to startread it →on this pageshowhide
explore
- Broker JMX Metrics and Health Signals6 questions
- Producer Throughput vs Latency Tuning6 questions
- Consumer Fetch and Parallelism Tuning5 questions
- Broker Threads, Page Cache and Disk Tuning5 questions
- Benchmarking with Perf Test Tools6 questions
- Client Metrics Instrumentation and Reporters5 questions
questions
page 2 of 2How does quota throttling show up in a Kafka request's latency breakdown, and how does the broker apply the throttle?
basics
~10 sWhen a client exceeds its configured quota, the broker delays the request's response and records that delay as ThrottleTimeMs. The delay is also sent back to the client so it can slow itself down.
One partition's leader broker shows much higher p99 than its peers even though cluster-wide throughput looks balanced. How would you diagnose and confirm a hot partition?
basics
~20 sA hot partition gets a disproportionate share of traffic, so its leader broker works harder and its tail latency rises while cluster averages look fine. Confirm by comparing per-partition/per-broker byte and message rates and checking the producer's partitioning key for skew.
How do throughput targets and metrics drive how many partitions a topic should have?
basics
~20 sPartition count must be at least target throughput divided by the per-partition throughput a single producer/consumer can sustain. Since one partition maps to one consumer in a group, partitions also set the max parallelism for consumers.
What do RequestHandlerAvgIdlePercent and NetworkProcessorAvgIdlePercent measure, and how do you interpret low values?
basics
~20 sBoth are saturation gauges from 0.0 to 1.0. RequestHandlerAvgIdlePercent is the fraction of time the I/O (request handler) threads are idle; NetworkProcessorAvgIdlePercent is the fraction of time the network threads are idle. Low values (near 0) mean those thread pools are saturated and the broker is a bottleneck.
What do IsrShrinksPerSec and IsrExpandsPerSec measure, and what does a high or flapping rate tell you?
basics
~20 sThey are rate meters counting how often replicas leave the in-sync replica set (shrink) or rejoin it (expand) per second. Occasional events are normal; a high or oscillating rate means replicas are repeatedly falling behind and catching up — a sign of an overloaded or unstable broker.
How should you size the JVM heap and choose a garbage collector for a Kafka broker, and why is a relatively small heap usually correct?
basics
~20 sGive the broker a modest heap (commonly ~5-6 GB) and leave most RAM to the OS page cache, since Kafka stores data off-heap in files. Use the G1 garbage collector (the default) to keep pauses low.
Compare JBOD vs RAID for Kafka brokers and explain why XFS is commonly recommended over ext4.
basics
~20 sJBOD gives Kafka each disk separately via log.dirs, so failures are isolated and replication handles redundancy; RAID (often RAID10) adds redundancy and even I/O but costs capacity/write speed. XFS is favored over ext4 for better performance under Kafka's large sequential-write workload.
Explain zero-copy / sendfile in Kafka: what it optimizes, when it applies, and what disables it.
basics
~20 sZero-copy uses the OS sendfile() call to stream bytes from the page cache straight to the network socket, skipping copies into and out of the JVM heap. It makes consumer fetches cheap — but only for plaintext data Kafka doesn't have to transform.
Explain KIP-714 client telemetry push: what problem it solves, how the protocol works, and how it relates to the MetricsReporter SPI.
basics
~20 sKIP-714 lets brokers collect standardized client metrics by having clients push them to the broker over the Kafka protocol, instead of operators scraping each client's JMX. The broker subscribes clients to metrics, clients push them as OpenTelemetry-encoded payloads on an interval.
A consumer group has growing lag but CPU on the instances is low and the network isn't saturated. Walk through how you'd diagnose and which fetch/parallelism configs you'd suspect.
basics
~20 sLow CPU + growing lag usually means the consumers are blocked waiting, not working: a slow downstream call, partition skew (one hot partition), too few partitions to parallelize, or fetches that are too small/infrequent. Check per-partition lag, then look at max.poll.records, fetch sizes, and partition assignment.
How does Burrow evaluate consumer-group health, and why is its sliding-window approach better than a single lag threshold?
basics
~20 sBurrow watches a sliding window of recent committed offsets per partition and looks at the trend, not one threshold. It flags a group as ERR/STALLED if offsets stop advancing while lag is nonzero, or WARN if lag is consistently growing, instead of alerting on an arbitrary fixed lag number.
Compare exposing consumer lag to Prometheus via kafka-exporter versus the JMX exporter. When would you use each?
basics
~20 skafka-exporter talks the Kafka protocol to the brokers and computes committed lag itself (it works without any live consumer). The JMX exporter scrapes MBeans from a JVM, so it surfaces client metrics like records-lag-max but only while that consumer is running. Use kafka-exporter for durable group lag; JMX exporter for in-process metrics.
What is buffer.memory, and what happens when a high-volume producer outpaces the brokers?
basics
~10 sbuffer.memory is the total bytes the producer can use to buffer unsent records. When it fills (the producer outpaces brokers), send() blocks up to max.block.ms, then throws TimeoutException. It's the producer's backpressure mechanism.
What does max.in.flight.requests.per.connection control, and how does it interact with retries, ordering, and idempotence?
basics
~20 sIt's the number of unacknowledged produce requests the producer allows per broker connection at once. Higher values pipeline more for throughput; with retries and no idempotence, values above 1 can reorder records on a failed-then-retried batch.
What is request purgatory in Kafka, and how does it relate to RemoteTimeMs for Fetch and Produce requests?
basics
~10 sPurgatory is where the broker parks requests that can't complete immediately and must wait for a condition (replication acks, or enough fetch data). That waiting time shows up as RemoteTimeMs.
A consumer reports high TotalTimeMs and you see RequestQueueTimeMs dominating with a large RequestQueueSize. What does this indicate and how do you remediate it?
basics
~10 sRequests 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.
A broker shows periodic p99 produce-latency spikes every 20-30 seconds that line up with LocalTimeMs jumps. How do you tell whether GC pauses or page-cache misses are the cause?
basics
~20 sPeriodic LocalTimeMs spikes point at the leader's local work stalling. Check GC pause logs/JMX for stop-the-world pauses lining up with the spikes; if GC is clean, look at page-cache misses forcing disk reads (rising disk read I/O and read latency while consumers fetch old data).
What types of client quotas does Kafka support and how do byte-rate versus request-rate quotas differ?
basics
~10 sKafka has two quota kinds: byte-rate quotas (producer_byte_rate, consumer_byte_rate) that cap MB/s per client, and request-rate quotas (request_percentage) that cap the share of broker request-handler/network thread time. They apply per client-id, user, or user+client-id.
How do you estimate the network bandwidth a Kafka cluster needs, including replication, from throughput metrics?
basics
~20 sOutbound bandwidth per broker ≈ consumer egress (BytesOutPerSec, multiplied by consumer fan-out) plus replication egress to followers (BytesInPerSec × (replicationFactor − 1)). Inbound ≈ producer ingress plus replica fetch traffic. Replication is often the dominant, easily-forgotten term.
What does LeaderElectionRateAndTimeMs measure, and how would you design alerting around broker health JMX metrics as a whole?
basics
~20 sLeaderElectionRateAndTimeMs is a timer the controller exposes: it tracks how often partition leader elections happen and how long they take (in ms). Frequent or slow elections signal instability. For overall broker health, alert on a small curated set — offline/under-replicated partitions, controller count, thread idle, ISR churn, election rate — with thresholds and dwell times tuned to severity.
Design producer configs for two workloads: (a) lowest p99 latency for small synchronous events, and (b) maximum throughput for a high-volume bulk ingest. Justify each setting.
basics
~10 sLow latency: linger.ms=0, small batches, compression none/lz4, acks=1, modest buffer. High throughput: linger.ms=20-100, large batch.size, zstd/lz4 compression, acks=all with pipelining, big buffer.memory. The two sit at opposite ends of the batching/latency curve.
Give me a repeatable, signal-driven decision procedure for isolating the bottleneck behind a Kafka p99/p999 latency spike across GC, page cache, slow followers, ISR shrink, request-queue saturation, network, hot partitions, and disk stalls.
basics
~20 sStart by splitting TotalTimeMs into phases to localize the stage, check idle-percent to spot pool saturation, then branch: RemoteTime to followers/ISR, LocalTime to GC/disk/page-cache, queue time to thread pools, send time to network. Compare per-broker to find a hot partition or one bad node.
What is Trogdor and when would you use it instead of the kafka-*-perf-test.sh scripts?
basics
~20 sTrogdor is Kafka's distributed test framework: a coordinator plus agents that run coordinated workloads and fault injections across many nodes. Use it for large-scale, multi-node, repeatable load and chaos tests; use the perf-test scripts for quick single-box benchmarks.
A producer shows high io-wait-ratio and low record-send-rate while CPU is idle. Walk through how you'd use client metrics to diagnose whether the bottleneck is the application, the producer config, or the broker/network.
basics
~20 sHigh io-wait-ratio with low send-rate means the producer's network thread is mostly idle waiting for data — it's not the bottleneck. Check record-queue-time, batch-size-avg, request-latency-avg, and buffer-available-bytes to see if the app, batching config, or broker is the limit.
showing 31–54 of 54