How does the MetricsReporter SPI work, and how would you use it to ship Kafka client metrics to an external system?
answer
- interface org.apache.kafka.common.metrics.MetricsReporter
- configure -> init -> metricChange/metricRemoval -> close
- register via metric.reporters config (CSV of class names)
- KafkaMetric.metricName().tags() + metricValue()
- in-process: must be non-blocking; JmxReporter is the default
basics
~20 sMetricsReporter 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).
solid answer
~40 sMetricsReporter is Kafka's SPI (service provider interface) for exporting client metrics without going through JMX. You implement `org.apache.kafka.common.metrics.MetricsReporter` — its lifecycle is `configure(map)`, then `init(List<KafkaMetric>)` with the initial metric set, then `metricChange(KafkaMetric)` and `metricRemoval(KafkaMetric)` callbacks as metrics appear/disappear, and `close()`. You register it by class name in the client config key `metric.reporters` (comma-separated; same config exists on brokers). Kafka instantiates it via reflection and passes all client configs to `configure`, so you can read custom keys. Each `KafkaMetric` exposes its `MetricName` (name, group, tags) and a live `metricValue()`. Common uses: a custom reporter that pushes to a TSDB, or established ones like `JmxReporter` (default), or wrapping a Prometheus/Micrometer bridge. It runs in-process, so reporters must be cheap and non-blocking to avoid stalling the client.
go deeper
Know that MetricsReporter is a pluggable interface registered via metric.reporters to export metrics outside JMX.
Implement the full lifecycle correctly and read MetricName tags + metricValue; know it's in-process.
Design a non-blocking, low-cardinality exporter and decide pull-snapshot vs push-callback semantics.
Standardize a reporter across services, manage cardinality/cost, and relate it to KIP-714's ClientTelemetryReporter.
**What it is.** An SPI (Service Provider Interface) is an extension point: a published interface you implement, and the framework discovers and calls your implementation. Kafka's metrics SPI is `org.apache.kafka.common.metrics.MetricsReporter`, and it lets you observe and export every metric a client produces, in-process, as an alternative or complement to JMX. **Lifecycle (the methods you implement):** - `void configure(Map<String, ?> configs)` — receives the full client configuration map (producer/consumer/admin configs). You can define and read your own keys here (e.g. an endpoint URL), since unknown configs are passed through. - `void init(List<KafkaMetric> metrics)` — the set of metrics that already exist at startup. - `void metricChange(KafkaMetric metric)` — called when a new metric is registered (and is the hook to start tracking it). Despite the name, registration of new metrics also arrives here. - `void metricRemoval(KafkaMetric metric)` — a metric is being unregistered (e.g. a per-topic metric for a topic no longer in use). - `void close()` — flush and release resources on client shutdown. - (Optionally `reconfigurable` configs / `ClusterResourceListener` if you need cluster id.) **KafkaMetric.** Each metric gives you a `MetricName` — which carries `name`, `group`, a description, and a `Map<String,String> tags` (client-id, topic, node-id, partition, etc.) — and `Object metricValue()` for the current value (a Double for most gauges/rates). You read `metricValue()` on your own schedule (pull) or react in the callbacks. **Registration.** Set the client config `metric.reporters` to a comma-separated list of fully-qualified class names. The default already includes `org.apache.kafka.common.metrics.JmxReporter` (which is what publishes the MBeans). Adding your reporter does not remove JMX unless you override the list. Brokers have the same `metric.reporters` config for server-side metrics. Related config: `metrics.sample.window.ms`, `metrics.num.samples`, `metrics.recording.level` (INFO vs DEBUG — DEBUG enables extra per-partition metrics). **How it's wired.** When a KafkaProducer/Consumer is constructed, it builds a `Metrics` object and instantiates each reporter class via reflection, calls `configure`, then `init`. From then on, as the client creates metrics lazily (e.g. first send to a topic creates that topic's metrics), `metricChange` fires. **Practical patterns & edge cases:** - **Keep it non-blocking.** Callbacks and any internal polling run alongside client threads; blocking I/O in `metricChange` can stall metric registration. Typical reporters buffer and flush on a background thread. - **Pull vs push.** Many reporters ignore the push callbacks for value updates and instead snapshot all known metrics periodically by calling `metricValue()`, because values change continuously but `metricChange` only fires on registration. - **Cardinality.** Per-topic/per-partition tags can explode cardinality in a TSDB; filter or aggregate. - **Don't confuse with broker-side `kafka.metrics.reporters`** (the old Yammer/KafkaMetricsReporter for server JMX) — the modern client/broker SPI is `metric.reporters` + `MetricsReporter`. - This same SPI is the foundation that KIP-714 builds on for the broker-pushed `ClientTelemetryReporter`.
- Your custom MetricsReporter does a synchronous HTTP POST inside metricChange. What's the risk?metricChange runs on the client's path when metrics register; a blocking/slow HTTP call stalls metric registration and can back up client operations. Buffer and flush asynchronously on a background thread instead; never block in the callback.
- If you set metric.reporters to only your reporter, what do you lose?You lose JmxReporter unless you re-add it — JMX MBeans disappear because the default list is overridden. Include org.apache.kafka.common.metrics.JmxReporter explicitly if you still want JMX.
saying these in an interview costs you the question
- Saying metrics are pushed to the reporter on every value change — metricChange fires on registration, not on each value update; you poll metricValue().
- Confusing the client SPI (metric.reporters / MetricsReporter) with the legacy broker kafka.metrics.reporters Yammer SPI.
- Doing blocking I/O directly in the callbacks.
- Assuming adding a reporter replaces JMX — it adds to the list unless you override it.