Kafka Producer Metrics in Java: What the Client Reports, How Micrometer Exposes It to Prometheus, and a Per-Record Latency Timer That Compiles

Lab run with apache/kafka:4.3.1 (KRaft), kafka-clients 4.3.1, Micrometer 1.17.1 and Java 17; 7 tests including javac runs of the earlier sample code

Revision note (2026-09-15). The earlier version of this article built its entire argument on a producer interceptor that read a header back through RecordMetadata.headers(). That method does not exist in any kafka-clients release; compiled against the versions the article named (kafka-clients 3.9.1, Micrometer 1.9.17) javac says cannot find symbol: method headers(). Against current versions the second sample fails as well, because micrometer-registry-prometheus moved packages in 1.13. The article also told readers to synchronize broker clocks for a measurement taken entirely on the producer JVM, promised "asynchronous and batch reporting" that Timer.record() does not do, and described expected /metrics output nobody had seen. It never named one of the metrics the producer already reports. This version comes from a lab: one Kafka 4.3.1 broker, kafka-clients 4.3.1, Micrometer 1.17.1, seven tests, output pasted as printed.

What was tested, and with what

  • Broker: apache/kafka:4.3.1 image, one KRaft node, plaintext, auto.create.topics.enable=false.
  • Client: org.apache.kafka:kafka-clients:4.3.1, io.micrometer:micrometer-core:1.17.1, io.micrometer:micrometer-registry-prometheus:1.17.1, Java 17, JUnit 5.13.4, Gradle 8.8.
  • Scrapes were done with Java's HttpClient against the lab's own endpoint. No Prometheus server, no Grafana.

Not tested: consumer and Kafka Streams metrics, interceptor overhead under sustained load, broker-side metrics, the Prometheus JMX exporter.

What a producer already reports

KafkaProducer.metrics() returns every metric the client maintains. After 100 records to one topic on one broker:

[lab] KafkaProducer.metrics(): 110 metrics in 5 groups
[lab]   group app-info (3 names)
[lab]   group kafka-metrics-count (1 names)
[lab]   group producer-metrics (73 names)
[lab]   group producer-node-metrics (12 names)
[lab]   group producer-topic-metrics (9 names)
[lab] record-send-total=100.0 request-latency-avg=14.0 record-queue-time-avg=105.0 version=4.3.1

The Kafka 4.3 monitoring documentation lists 22 names under "Producer Sender Metrics" (batch-size-avg, batch-size-max, batch-split-rate, batch-split-total, compression-rate-avg, metadata-age, produce-throttle-time-avg, produce-throttle-time-max, record-error-rate, record-error-total, record-queue-time-avg, record-queue-time-max, record-retry-rate, record-retry-total, record-send-rate, record-send-total, record-size-avg, record-size-max, records-per-request-avg, request-latency-avg, request-latency-max, requests-in-flight) and 9 per-topic names. The lab asserts that all of them are present in the producer-metrics and producer-topic-metrics groups of the 4.3.1 client. The other 51 names in producer-metrics are the common client metrics (connections, bytes in and out, I/O thread time, authentication, buffer pool, transactions) that the same documentation page describes in its "Common monitoring metrics" section.

Three of them answer most "is my producer healthy" questions on their own:

  • request-latency-avg / request-latency-max: time per produce request to the broker, in ms. Per request, not per record; one request carries a batch.
  • record-queue-time-avg / record-queue-time-max: how long records waited in the accumulator before being sent. In the run above this was 105 ms against a request latency of 14 ms, because the test sent 100 records in a tight loop and called flush() once; the wait is batching, not the network.
  • record-error-total: records that failed for good. Sends that never reach a broker count here too (below).

JMX has the same names

The JMX reporter is on by default and registers one MBean per group and tag combination under the kafka.producer domain:

[lab] JMX MBeans for client-id=lab-jmx: 6
[lab]   kafka.producer:type=app-info,client-id=lab-jmx
[lab]   kafka.producer:type=kafka-metrics-count,client-id=lab-jmx
[lab]   kafka.producer:type=producer-metrics,client-id=lab-jmx
[lab]   kafka.producer:type=producer-node-metrics,client-id=lab-jmx,node-id=node--1
[lab]   kafka.producer:type=producer-node-metrics,client-id=lab-jmx,node-id=node-1
[lab]   kafka.producer:type=producer-topic-metrics,client-id=lab-jmx,topic=metrics-f239717f
[lab] producer-metrics attribute count: JMX=73 metrics()=73 equal=true
[lab] JMX record-send-total=25.0

The 73 attribute names on the producer-metrics MBean are exactly the 73 names from metrics() (asserted). node--1 is the bootstrap connection before the client learned the real node id. If you run the Prometheus JMX exporter, these are the MBeans it sees; the rest of this article does it in-process instead.

Exposing the metrics to Prometheus with two dependencies and no custom code

Micrometer ships a binder, io.micrometer.core.instrument.binder.kafka.KafkaClientMetrics, that reads producer.metrics() and registers the result as Micrometer meters. With the Prometheus registry, the exposition text comes from registry.scrape(); serving it needs nothing beyond the JDK.

dependencies {
    implementation 'org.apache.kafka:kafka-clients:4.3.1'
    implementation 'io.micrometer:micrometer-core:1.17.1'
    implementation 'io.micrometer:micrometer-registry-prometheus:1.17.1'
}
import com.sun.net.httpserver.HttpServer;
import io.micrometer.core.instrument.binder.kafka.KafkaClientMetrics;
import io.micrometer.prometheusmetrics.PrometheusConfig;
import io.micrometer.prometheusmetrics.PrometheusMeterRegistry;

public final class PrometheusExposure implements AutoCloseable {
    private final PrometheusMeterRegistry registry = new PrometheusMeterRegistry(PrometheusConfig.DEFAULT);
    private final KafkaClientMetrics binder;
    private HttpServer server;

    public PrometheusExposure(Producer<?, ?> producer) {
        this.binder = new KafkaClientMetrics(producer);
        this.binder.bindTo(registry);
    }

    public void serve(int port) throws IOException {
        server = HttpServer.create(new InetSocketAddress("127.0.0.1", port), 0);
        server.createContext("/metrics", exchange -> {
            byte[] body = registry.scrape().getBytes(StandardCharsets.UTF_8);
            exchange.getResponseHeaders().add("Content-Type", "text/plain; version=0.0.4; charset=utf-8");
            exchange.sendResponseHeaders(200, body.length);
            try (OutputStream out = exchange.getResponseBody()) {
                out.write(body);
            }
        });
        server.start();
    }

    @Override
    public void close() {
        if (server != null) {
            server.stop(0);
        }
        binder.close();
        registry.close();
    }
}

The binder holds a scheduler that refreshes the meter set, so close() matters when the producer goes away. Scraped after 40 records:

[lab] GET /metrics -> 200, 256 lines, 86 metric families
[lab] Micrometer meter names bound by KafkaClientMetrics: 86
[lab]   kafka_app_info_start_time_ms{client_id="lab-prom",kafka_version="4.3.1"} 1.789521303175E12
[lab]   kafka_producer_record_error_total{client_id="lab-prom",kafka_version="4.3.1"} 0.0
[lab]   kafka_producer_record_queue_time_avg{client_id="lab-prom",kafka_version="4.3.1"} 1.0
[lab]   kafka_producer_record_queue_time_max{client_id="lab-prom",kafka_version="4.3.1"} 1.0
[lab]   kafka_producer_record_send_rate{client_id="lab-prom",kafka_version="4.3.1"} 1.32687587076229
[lab]   kafka_producer_record_send_total{client_id="lab-prom",kafka_version="4.3.1"} 40.0
[lab]   kafka_producer_request_latency_avg{client_id="lab-prom",kafka_version="4.3.1"} 6.0
[lab]   kafka_producer_request_latency_max{client_id="lab-prom",kafka_version="4.3.1"} 6.0

Naming rule, as observed: Kafka group producer-metrics plus name record-send-total becomes the Micrometer meter kafka.producer.record.send.total, which the Prometheus registry renders as kafka_producer_record_send_total, tagged with client_id and kafka_version (plus topic or node_id for the per-topic and per-node groups). The 110 Kafka metrics collapse to 86 meter names because the per-topic and per-node groups share names across tag values.

Why the old sample no longer compiles

Micrometer 1.13 rebuilt micrometer-registry-prometheus on the Prometheus Java client 1.x (io.prometheus:prometheus-metrics-core) and moved the classes to io.micrometer.prometheusmetrics. The previous package, io.micrometer.prometheus, and its CollectorRegistry-based API now live in io.micrometer:micrometer-registry-prometheus-simpleclient. Code written for 1.9 that also pulls io.prometheus.client.exporter.HTTPServer from simpleclient_httpserver therefore fails against 1.17.1 with:

KafkaProducerWithMicrometer.java:7: package io.micrometer.prometheus does not exist
KafkaProducerWithMicrometer.java:9: package io.prometheus.client.exporter does not exist

Either switch to the prometheusmetrics package and drop the simpleclient server (as above), or depend on the -simpleclient registry artifact and keep the old imports. Mixing the two is what produces the nine errors the lab recorded.

A per-record latency timer that compiles

The built-in request-latency-* metrics are per produce request. If you want the time from send() to acknowledgement per record, an interceptor can measure it, but not the way the earlier article did. ProducerInterceptor.onAcknowledgement(RecordMetadata, Exception) receives no headers, and RecordMetadata has no headers() method; its accessors are offset, timestamp, topic, partition, serializedKeySize, serializedValueSize.

Since kafka-clients 4.1.0, ProducerInterceptor has a second default method, onAcknowledgement(RecordMetadata metadata, Exception exception, Headers headers), which receives the record's read-only headers (checked in the 4.0.0 and 4.1.0 sources jars: absent in the first, present in the second). Override that one:

public final class SendLatencyInterceptor implements ProducerInterceptor<String, String> {
    public static final String REGISTRY_CONFIG = "meter.registry";
    static final String HEADER = "send-timestamp-nanos";
    private Timer timer;

    @Override
    public void configure(Map<String, ?> configs) {
        Object registry = configs.get(REGISTRY_CONFIG);
        if (!(registry instanceof MeterRegistry)) {
            throw new IllegalArgumentException(REGISTRY_CONFIG + " must be a MeterRegistry, got " + registry);
        }
        this.timer = Timer.builder("kafka.producer.record.latency")
                .description("send() to acknowledgement, measured on the producer's clock")
                .register((MeterRegistry) registry);
    }

    @Override
    public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
        record.headers().add(HEADER, Long.toString(System.nanoTime()).getBytes(StandardCharsets.UTF_8));
        return record;
    }

    @Override
    public void onAcknowledgement(RecordMetadata metadata, Exception exception, Headers headers) {
        if (exception != null || headers == null) {
            return;
        }
        Header header = headers.lastHeader(HEADER);
        if (header == null) {
            return;
        }
        long sentAt = Long.parseLong(new String(header.value(), StandardCharsets.UTF_8));
        timer.record(System.nanoTime() - sentAt, TimeUnit.NANOSECONDS);
    }

    @Override
    public void close() {
    }
}

Wiring: put the class name under interceptor.classes and the registry object itself under meter.registry in the producer Properties; KafkaProducer passes non-string config values through to configure() unchanged (this is what the lab does, and it works). Fifty records:

[lab] interceptor timer: count=50 mean=7.45 ms max=13.45 ms
[lab] onAcknowledgement ran on thread: kafka-producer-network-thread | lab-interceptor
[lab] built-in request-latency-avg for the same producer: 4.0 ms (per produce request, not per record)

Three things the output settles:

  • The timer counted exactly the 50 records sent (asserted), so the header survives the trip through the accumulator and comes back with the acknowledgement.
  • Both timestamps come from System.nanoTime() on the producer JVM. Broker clock skew cannot enter this number; the earlier article's advice to synchronize broker clocks for it was wrong. (Clock skew does matter for a different measurement: consumer-side end-to-end latency computed from the record timestamp.)
  • The callback runs on kafka-producer-network-thread, the single I/O thread of the producer. Timer.record() is a synchronous in-memory update and is fine there; anything that blocks or allocates heavily is not, because every other in-flight record waits behind it. There is no asynchronous or batched reporting inside Micrometer's Timer to fall back on.

The header adds about 20 bytes per record and is visible to consumers. If that is unacceptable, the alternative is a map from a record identity to a start time inside the interceptor, with its own eviction problem; the lab did not build that.

Failures show up in record-error-total, even when nothing reached the broker

With auto.create.topics.enable=false and max.block.ms=500, three sends to a topic that does not exist:

[lab] send() to a missing topic with auto.create.topics.enable=false, max.block.ms=500 -> TimeoutException: Topic does-not-exist-ec70578a not present in metadata after 500 ms.
[lab] record-error-total=3.0 record-send-total=0.0

The interceptor above would have seen these as onAcknowledgement(null-ish metadata, TimeoutException, headers) and skipped them, which is why a separate error count is worth having. kafka_producer_record_error_total is on the scrape without further work.

Sources

  • Kafka 4.3 documentation, Monitoring (producer sender metrics and common client metrics): https://kafka.apache.org/43/operations/monitoring/
  • Kafka 4.3.1 javadoc, ProducerInterceptor (both onAcknowledgement overloads): https://kafka.apache.org/43/javadoc/org/apache/kafka/clients/producer/ProducerInterceptor.html
  • Kafka 4.3.1 javadoc, RecordMetadata: https://kafka.apache.org/43/javadoc/org/apache/kafka/clients/producer/RecordMetadata.html
  • Micrometer reference, Kafka instrumentation (KafkaClientMetrics): https://docs.micrometer.io/micrometer/reference/reference/kafka.html
  • Micrometer reference, Prometheus registry: https://docs.micrometer.io/micrometer/reference/implementations/prometheus.html
  • Micrometer 1.13 migration guide (Prometheus client 1.x, -simpleclient artifact): https://github.com/micrometer-metrics/micrometer/wiki/1.13-Migration-Guide
  • Prometheus JMX exporter, for the JMX route instead: https://github.com/prometheus/jmx_exporter
  • kafka-clients 4.3.1 on Maven Central: https://central.sonatype.com/artifact/org.apache.kafka/kafka-clients/4.3.1