如何在Scala中通过Micrometer暴露Kafka客户端内部指标
在Scala中暴露Kafka客户端内置指标入门指引
核心思路
Kafka Java客户端自带Metrics系统,所有内置指标都通过KafkaConsumer.metrics()方法暴露。我们可以通过以下步骤提取目标指标并对外暴露:
- 引入必要依赖(Kafka客户端、指标收集库)
- 初始化KafkaConsumer并获取Metrics实例
- 定位目标指标的标识(名称+标签)
- 将Kafka指标适配到通用指标库(如Micrometer)
- 启动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
相关产品推荐
相关产品推荐

