You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Scala中通过Micrometer暴露Kafka客户端内部指标

在Scala中暴露Kafka客户端内置指标入门指引

核心思路

Kafka Java客户端自带Metrics系统,所有内置指标都通过KafkaConsumer.metrics()方法暴露。我们可以通过以下步骤提取目标指标并对外暴露:

  1. 引入必要依赖(Kafka客户端、指标收集库)
  2. 初始化KafkaConsumer并获取Metrics实例
  3. 定位目标指标的标识(名称+标签)
  4. 将Kafka指标适配到通用指标库(如Micrometer)
  5. 启动HTTP端点暴露指标(供Prometheus等采集)

步骤1:添加依赖(SBT示例)

在build.sbt中加入以下依赖:

libraryDependencies ++= Seq(
  "org.apache.kafka" % "kafka-clients" % "3.6.0", // 匹配你的Kafka集群版本
  "io.micrometer" % "micrometer-core" % "1.11.0",
  "io.micrometer" % "micrometer-registry-prometheus" % "1.11.0",
  "org.eclipse.jetty" % "jetty-server" % "11.0.15" // 用于启动简单HTTP服务暴露指标
)

步骤2:提取并适配Kafka指标

以下是完整Scala示例代码,包含消费者初始化、指标提取、暴露逻辑:

import org.apache.kafka.clients.consumer.{KafkaConsumer, ConsumerConfig, ConsumerRecords}
import org.apache.kafka.common.serialization.StringDeserializer
import io.micrometer.core.instrument.binder.kafka.KafkaConsumerMetrics
import io.micrometer.prometheus.PrometheusMeterRegistry
import org.eclipse.jetty.server.Server
import org.eclipse.jetty.servlet.{ServletContextHandler, ServletHolder}
import io.prometheus.client.exporter.MetricsServlet

import java.time.Duration
import java.util.Properties
import scala.jdk.CollectionConverters._

object KafkaMetricsExporter {
  def main(args: Array[String]): Unit = {
    // 初始化Kafka消费者配置
    val consumerProps = new Properties()
    consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
    consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "test-metrics-group")
    consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer].getName)
    consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer].getName)
    consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")

    // 创建KafkaConsumer实例
    val consumer = new KafkaConsumer[String, String](consumerProps)
    consumer.subscribe(List("test-topic").asJava)

    // 绑定Kafka消费者指标到Micrometer的Prometheus注册表
    val meterRegistry = new PrometheusMeterRegistry(PrometheusMeterRegistry.DEFAULT_CONFIG)
    new KafkaConsumerMetrics(consumer).bindTo(meterRegistry)

    // 启动Jetty服务器暴露Prometheus指标端点
    val server = new Server(8080)
    val context = new ServletContextHandler()
    context.setContextPath("/")
    server.setHandler(context)
    context.addServlet(new ServletHolder(new MetricsServlet(meterRegistry.getPrometheusRegistry)), "/metrics")
    server.start()
    println("Metrics endpoint started at http://localhost:8080/metrics")

    // 模拟消费逻辑(保持消费者运行以生成指标)
    while (true) {
      val records: ConsumerRecords[String, String] = consumer.poll(Duration.ofMillis(100))
      // 处理消息(实际业务逻辑替换此处)
      records.asScala.foreach(record => println(s"Consumed record: ${record.value()}"))
    }
  }
}

步骤3:验证目标指标

启动程序后,访问http://localhost:8080/metrics,可以找到以下对应指标:

  • kafka_consumer_commit_time_total_sum:对应committed-time-ns-total(累计提交耗时,单位纳秒)
  • kafka_consumer_fetch_size_avg:对应fetch-size-avg(平均拉取大小,单位字节)
  • kafka_consumer_records_consumed_rate:对应records-consumed-rate(每秒消费记录数)

关键说明

  • 若需确认Kafka指标的准确标识,可通过consumer.metrics().keySet()打印所有指标的名称和标签
  • Micrometer的KafkaConsumerMetrics已封装大部分常用指标的适配,无需手动逐个提取
  • 若不想依赖Micrometer,也可直接遍历consumer.metrics()获取指标值,自行实现HTTP暴露逻辑

内容的提问来源于stack exchange,提问作者Vimal

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 03:42:52