skip to content

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 pageshow

questions

page 1 of 2

What is kafka-producer-perf-test.sh, and what do its --num-records, --record-size, and --throughput flags control?

level: juniorimportance: must knowfreq 62%

answer

  1. num-records = total count, stops the run
  2. record-size = bytes per message → drives MB/s
  3. throughput = records/sec cap; -1 = flat out
  4. must add --topic + --producer-props bootstrap.servers
  5. reports records/sec, MB/sec, p99 latency

basics

~10 s

It is Kafka's built-in load generator for producers. --num-records sets how many messages to send, --record-size sets each message's size in bytes, and --throughput caps messages per second (-1 means unthrottled, full speed).

solid answer

~40 s

kafka-producer-perf-test.sh is a CLI tool shipped with Kafka that generates synthetic producer load and reports throughput and latency. --num-records is the total count of records to publish (the test stops after that many). --record-size is the per-record payload size in bytes, so total data sent ≈ num-records × record-size. --throughput is the target rate in records/second; set it to a positive number to throttle to a fixed rate (useful for latency tests at a controlled load) or to -1 to send as fast as possible (max-throughput tests). You also pass --topic and --producer-props (e.g. bootstrap.servers, acks, batch.size, linger.ms, compression.type). Output reports records/sec, MB/sec, and latency percentiles (avg, p50, p95, p99, max).

go deeper

for a junior

Know the three flags by name and that throughput is records/sec with -1 meaning unthrottled.

for a middle

Connect record-size to MB/s and know the required --topic/--producer-props and what the summary reports.

for a senior

Discuss payload compressibility, parallel producer instances, and choosing throttle vs flat-out per measurement goal.

for a principal

Frame the tool's role in a repeatable benchmark methodology and its limits as a synthetic single-JVM generator.

## What the tool is `kafka-producer-perf-test.sh` is a command-line benchmarking utility bundled in every Kafka distribution's `bin/` directory (the wrapper around the `org.apache.kafka.tools.ProducerPerformance` class). Its job is to push a controlled stream of synthetic records into a topic and measure how fast the producer can go and how long each send takes, without you having to write any code. ## The three core flags - **`--num-records`**: the total number of records the test will send before stopping. This bounds the run. Bigger values give more statistically stable numbers and let the system reach steady state, but take longer. - **`--record-size`**: the size, in bytes, of each record's value payload. The tool fills records with random bytes of this size. Total bytes pushed ≈ `num-records × record-size`. This matters because Kafka throughput is often network/disk-bandwidth bound, so MB/s depends directly on record size. (Alternatively `--payload-file` feeds real payloads.) - **`--throughput`**: the target send rate in **records per second**. A positive value throttles the producer to that rate using a rate limiter — essential when you want to measure *latency at a fixed offered load*. Setting it to `-1` removes the throttle so the producer runs flat-out, which is how you measure *maximum sustainable throughput*. ## Required companions You must also pass `--topic <name>` and `--producer-props key=value ...` (at minimum `bootstrap.servers`). The producer-props let you vary the knobs that actually move throughput/latency: `acks`, `batch.size`, `linger.ms`, `compression.type`, `buffer.memory`. There is also `--print-metrics` to dump the full client metric set at the end. ## Output The tool prints periodic progress lines and a final summary like: `1000000 records sent, 250000.0 records/sec (23.84 MB/sec), 5.20 ms avg latency, 120.00 ms max latency, ... 99th 45 ms`. Throughput is reported in both **records/sec** and **MB/sec**; latency is reported as average, max, and percentiles (p50/p95/p99). ## Edge cases / gotchas - Random payloads compress poorly, so if you enable `compression.type` your MB/s on the wire won't reflect real (compressible) data — use `--payload-file` with representative data. - A single producer process is single-JVM; one instance may not saturate a large cluster, so you may run several in parallel. - The reported MB/sec is uncompressed payload throughput from the client's perspective.

  • How do you compute the total data volume a run will push?
    Approximately num-records × record-size bytes of payload (e.g. 1,000,000 records × 1000 bytes ≈ 1 GB), before any compression and excluding Kafka record/batch overhead.
  • When would you set --throughput to -1 versus a fixed number?
    Use -1 to find maximum sustainable throughput (unthrottled). Use a fixed positive value to hold a controlled offered load while measuring latency, which is how you build a throughput-vs-latency curve.

saying these in an interview costs you the question

  • Saying --throughput is in MB/s — it is records/second.
  • Thinking the tool needs custom code; it's a ready CLI in bin/.
  • Ignoring that random payloads make compression results meaningless.
  • Believing one producer instance always saturates the cluster.

context

open as a page

What does the UnderReplicatedPartitions broker metric mean, and what should it normally read?

level: juniorimportance: must knowfreq 70%

basics

~20 s

UnderReplicatedPartitions counts how many partitions led by this broker have fewer in-sync replicas than the configured replication factor. In a healthy cluster it should be 0. A sustained non-zero value means replicas are falling behind or down.

open as a page

What is JMX in the context of a Kafka broker, and how do you expose and read broker metrics through it?

level: juniorimportance: must knowfreq 55%

basics

~20 s

JMX (Java Management Extensions) is the Java standard for exposing runtime metrics as MBeans. A Kafka broker publishes its internal metrics as JMX MBeans. You enable it by setting JMX_PORT (or jmxremote system properties) and read it with tools like JConsole, jmxterm, or a Prometheus JMX exporter.

open as a page

What do num.network.threads and num.io.threads control on a Kafka broker, and how do you decide how many to set?

level: juniorimportance: must knowfreq 70%

basics

~10 s

num.network.threads handle reading requests off and writing responses onto network sockets; num.io.threads do the actual work (reading/writing the log on disk). Default 3 and 8. Raise them if request-handler/network idle time drops.

open as a page

What are the most important built-in JMX metrics for a Kafka producer and consumer, and what does each tell you about client health?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Kafka clients expose metrics over JMX. Key producer metrics: record-send-rate, request-latency-avg, buffer-available-bytes. Key consumer metrics: records-consumed-rate, fetch-latency-avg, records-lag-max. They show throughput, latency, and buffering/lag health.

open as a page

What is max.poll.records and why might you lower it when a consumer keeps getting kicked out of its group?

level: juniorimportance: must knowfreq 70%

basics

~20 s

max.poll.records caps how many records one poll() call returns. If processing each batch takes too long, the consumer misses its deadline and gets removed from the group. Lowering it means smaller batches that process faster, keeping the consumer alive.

open as a page

What is consumer lag in Kafka, and how do you compute it for a single partition?

level: juniorimportance: must knowfreq 80%

basics

~20 s

Consumer lag is how far behind a consumer is on a partition: log-end offset minus the consumer's committed offset. A lag of 0 means the consumer has read everything; a growing lag means it can't keep up with producers.

open as a page

What do batch.size and linger.ms do on a Kafka producer, and how do they interact to trade latency for throughput?

level: juniorimportance: must knowfreq 78%

basics

~10 s

batch.size is the max bytes per partition batch; linger.ms is how long the producer waits to fill a batch before sending. Raising either groups more records per request, boosting throughput but adding latency.

open as a page

Which broker-side metric tells you the end-to-end time a Kafka request took, and where do you find it?

level: juniorimportance: must knowfreq 60%

basics

~10 s

TotalTimeMs, exposed via JMX under kafka.network:type=RequestMetrics, name=TotalTimeMs, request=<ApiKey> (like Produce or Fetch). It's a histogram, so you read Mean, 99thPercentile, and Max.

open as a page

What do the BytesInPerSec, BytesOutPerSec, and MessagesInPerSec broker metrics measure, and where do you find them?

level: juniorimportance: must knowfreq 70%

basics

~20 s

They are JMX rate meters on the broker: BytesInPerSec is bytes produced into the broker per second, BytesOutPerSec is bytes consumers fetch out per second, and MessagesInPerSec is records produced per second. They exist both per-broker and per-topic.

open as a page

What do ActiveControllerCount and OfflinePartitionsCount tell you, and what is healthy for each?

level: middleimportance: must knowfreq 60%

basics

~20 s

ActiveControllerCount is per-broker and is 1 on the single controller broker and 0 on the rest, so cluster-wide it must sum to exactly 1. OfflinePartitionsCount is the number of partitions with no leader — healthy is 0; any positive value means those partitions are unavailable.

open as a page

Why does Kafka rely on the OS page cache instead of flushing every write to disk, and what is the role of log.flush.interval.messages / log.flush.interval.ms?

level: middleimportance: must knowfreq 65%

basics

~20 s

Kafka writes to the OS page cache and lets the OS flush to disk in the background, which is fast. Durability comes from replication, not fsync. log.flush.interval.* would force periodic fsyncs, but Kafka leaves them effectively off by default.

open as a page

Explain how fetch.min.bytes and fetch.max.wait.ms work together to trade latency for throughput on the consumer.

level: middleimportance: must knowfreq 65%

basics

~20 s

fetch.min.bytes tells the broker not to answer a fetch until it has at least that many bytes ready. fetch.max.wait.ms caps how long it waits for that. Bigger min.bytes = fewer, larger fetches (more throughput, more latency); the wait prevents waiting forever when traffic is low.

open as a page

How do you read kafka-consumer-groups --describe output, and why is offset lag not the same as end-to-end latency?

level: middleimportance: must knowfreq 60%

basics

~20 s

kafka-consumer-groups --describe shows, per partition, CURRENT-OFFSET (committed), LOG-END-OFFSET, LAG (the difference), plus consumer-id/host. Offset lag counts messages behind. End-to-end latency is the time between when a record was produced and when it is consumed — a time, not a count — so high-throughput partitions can have big offset lag but tiny time latency.

open as a page

What are the records-lag and records-lag-max client metrics, and why might they differ from what kafka-consumer-groups reports?

level: middleimportance: must knowfreq 65%

basics

~20 s

records-lag is a per-partition consumer client (JMX) metric showing how far behind that fetch is; records-lag-max is the maximum across the consumer's assigned partitions. They come from the live client, so they only exist while the consumer is fetching and reflect what it has fetched, not committed.

open as a page

How does the acks setting trade durability against producer throughput and latency?

level: middleimportance: must knowfreq 80%

basics

~20 s

acks controls how many brokers must confirm a write. acks=0 (no wait) is fastest but can lose data; acks=1 waits for the leader; acks=all waits for all in-sync replicas — most durable but highest latency.

open as a page

What are the phases of a Kafka broker request as exposed by RequestMetrics, and how do they sum into TotalTimeMs?

level: middleimportance: must knowfreq 70%

basics

~20 s

Kafka splits each request's time into phases: time waiting in the request queue, time processing locally, time waiting on other brokers, time in the response queue, and time sending the response. TotalTimeMs is their sum.

open as a page

What do the network-processor and request-handler idle-percent metrics tell you, and how do you use them to spot thread-pool saturation behind tail-latency spikes?

level: middleimportance: must knowfreq 70%

basics

~20 s

Idle percent is the fraction of time a thread pool sits with no work. NetworkProcessorAvgIdlePercent and RequestHandlerAvgIdlePercent near 0 mean the network or I/O threads are saturated, which causes requests to queue up and tail latency to spike.

open as a page

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%

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.

open as a page

How do you estimate disk capacity for a Kafka cluster from throughput metrics and retention settings?

level: middleimportance: must knowfreq 60%

basics

~10 s

Disk needed = ingress bytes/sec × retention seconds × replication factor, summed across topics, plus headroom. Use BytesInPerSec for the rate and retention.ms (or retention.bytes) for how long data is kept.

open as a page

How do you use the perf-test tools to build a throughput-vs-latency curve, and how do you interpret its 'knee'?

level: seniorimportance: must knowfreq 40%

basics

~20 s

Run kafka-producer-perf-test.sh repeatedly at increasing fixed --throughput values and record p99 latency each time. Plot offered rate (x) against latency (y). Latency stays flat then sharply rises at the 'knee' — the point where you hit saturation.

open as a page

Why does adding more consumer instances to a group eventually stop improving throughput, and how do partitions set that ceiling?

level: seniorimportance: must knowfreq 68%

basics

~20 s

Within a consumer group, each partition is consumed by exactly one member. So the number of partitions is the maximum useful parallelism. Once you have as many consumers as partitions, extra consumers sit idle and throughput stops rising.

open as a page

Producers using acks=all see p99 latency spikes, and you notice ISR shrinking on some partitions. Explain the chain of causation and which metrics you'd correlate to confirm a slow follower is the culprit.

level: seniorimportance: must knowfreq 62%

basics

~20 s

With acks=all a produce request can't complete until all in-sync replicas (the ISR) catch up. A slow follower lags, so the produce waits longer (high RemoteTimeMs); if it lags past replica.lag.time.max.ms the leader drops it from ISR (ISR shrink). Correlate RemoteTimeMs, replica lag, ISR-shrink rate, and the follower's own health.

open as a page

How does kafka-consumer-perf-test.sh measure consumer performance, and how do you interpret its output?

level: middleimportance: should knowfreq 45%

basics

~20 s

It runs a consumer that reads a fixed number of messages from a topic and reports how much data and how many records it consumed per second. You read its MB/sec and nMsg/sec columns to judge consume throughput.

open as a page

When benchmarking Kafka, what does it mean to measure 'sustained MB/s' versus 'records/s', and why can they tell different stories?

level: middleimportance: should knowfreq 30%

basics

~20 s

MB/s is data-volume throughput (bytes per second); records/s is message-count throughput (messages per second). They differ because record size links them: small records can give high records/s but low MB/s, and vice versa. 'Sustained' means the rate held over a long, steady run, not a peak burst.

open as a page

Why must you account for warm-up effects when benchmarking Kafka, and how do you isolate steady-state throughput?

level: middleimportance: should knowfreq 33%

basics

~20 s

Early in a run the JVM JIT, OS page cache, TCP connections, and producer batching haven't stabilized, so the first numbers are misleadingly slow. Run long enough and discard the warm-up interval, quoting only the steady-state rate.

open as a page

How do you integrate Kafka client metrics into a Micrometer/OpenTelemetry-based observability stack (e.g. a Spring Boot service exporting to Prometheus)?

level: middleimportance: should knowfreq 40%

basics

~20 s

Bind Kafka's client metrics into Micrometer using KafkaClientMetrics (or Spring Boot's auto-binding for KafkaTemplate/listener containers). Micrometer then exports them to your backend (Prometheus, OTel). It reads the client's metrics() map and registers them as Micrometer gauges.

open as a page

How does the MetricsReporter SPI work, and how would you use it to ship Kafka client metrics to an external system?

level: middleimportance: should knowfreq 45%

basics

~20 s

MetricsReporter is a pluggable interface (org.apache.kafka.common.metrics.MetricsReporter). You implement it, register it via the metric.reporters client config, and Kafka calls your code as metrics are created/changed/removed so you can forward them (e.g. to Prometheus or a custom sink).

open as a page

What is max.partition.fetch.bytes versus fetch.max.bytes, and what breaks if max.partition.fetch.bytes is smaller than your largest message?

level: middleimportance: should knowfreq 50%

basics

~20 s

fetch.max.bytes caps the total size of a whole fetch response across all partitions; max.partition.fetch.bytes caps how much comes from each single partition. Historically, if a message was bigger than max.partition.fetch.bytes the consumer could stall — modern clients still return that oversized record so the consumer can progress.

open as a page

How does compression.type affect producer throughput, and how do lz4, zstd, and snappy compare?

level: middleimportance: should knowfreq 62%

basics

~10 s

compression.type compresses each batch before sending, cutting network and disk usage and raising throughput. snappy/lz4 are fast with modest ratios; zstd compresses harder (smaller payloads) at more CPU; gzip is highest ratio but slowest.

open as a page

showing 1–30 of 54