如何在非Spring Boot的核心Java应用中暴露Prometheus API并实现Micrometer
核心Java流式应用集成Micrometer实现指南
1. 添加依赖
在项目中引入Micrometer核心库与Prometheus导出器(用于对接Prometheus监控系统),以Maven为例:
<dependencies> <!-- Micrometer Core --> <dependency> <groupId>io.micrometer</groupId> <artifactId>micrometer-core</artifactId> <version>1.12.0</version> </dependency> <!-- Micrometer Prometheus Exporter --> <dependency> <groupId>io.micrometer</groupId> <artifactId>micrometer-registry-prometheus</artifactId> <version>1.12.0</version> </dependency> </dependencies>
2. 初始化全局MeterRegistry
创建单例模式的PrometheusMeterRegistry,避免重复实例化,同时配置通用标签区分实例/环境:
import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.config.MeterFilter; import io.micrometer.prometheus.PrometheusMeterRegistry; public class MetricsConfig { private static final PrometheusMeterRegistry registry; static { registry = new PrometheusMeterRegistry(PrometheusConfig.DEFAULT); // 添加通用标签,用于多实例/多环境区分 registry.config().commonTags("application", "stream-processing-app", "env", "production"); // 可选:过滤不需要的指标(如部分JVM GC指标) registry.config().meterFilter(MeterFilter.denyNameStartsWith("jvm.gc")); } public static MeterRegistry getRegistry() { return registry; } public static String scrapeMetrics() { return registry.scrape(); } }
3. 流式场景指标埋点
针对流式处理的核心场景,添加对应的监控指标:
3.1 消息处理成功/失败计数
用Counter统计不同结果的消息数量:
import io.micrometer.core.instrument.Counter; import io.micrometer.core.instrument.MeterRegistry; public class StreamProcessor { private final Counter successMessages; private final Counter failedMessages; public StreamProcessor() { MeterRegistry registry = MetricsConfig.getRegistry(); successMessages = Counter.builder("stream.messages.processed") .tag("result", "success") .description("Total number of successfully processed messages") .register(registry); failedMessages = Counter.builder("stream.messages.processed") .tag("result", "failed") .description("Total number of failed messages") .register(registry); } public void processMessage(String message) { try { // 模拟消息处理逻辑 Thread.sleep(10); successMessages.increment(); } catch (Exception e) { failedMessages.increment(); } } }
3.2 消息处理耗时统计
用Timer记录单条消息的处理时长:
import io.micrometer.core.instrument.Timer; public class StreamProcessor { // 省略之前的计数器定义 private final Timer processingTimer; public StreamProcessor() { MeterRegistry registry = MetricsConfig.getRegistry(); // 省略计数器初始化 processingTimer = Timer.builder("stream.message.processing.duration") .description("Duration of message processing") .register(registry); } public void processMessage(String message) { processingTimer.record(() -> { try { // 模拟消息处理逻辑 Thread.sleep(10); successMessages.increment(); } catch (Exception e) { failedMessages.increment(); } }); } }
3.3 待处理队列长度监控
如果应用使用队列缓存消息,用Gauge实时监控队列大小:
import io.micrometer.core.instrument.Gauge; import java.util.Queue; import java.util.concurrent.ConcurrentLinkedQueue; public class MessageQueue { private final Queue<String> queue = new ConcurrentLinkedQueue<>(); public MessageQueue() { MeterRegistry registry = MetricsConfig.getRegistry(); Gauge.builder("stream.queue.size", queue, Queue::size) .description("Number of pending messages in queue") .register(registry); } public void addMessage(String message) { queue.add(message); } public String takeMessage() { return queue.poll(); } }
4. 暴露指标端点
核心Java无内置HTTP端点,需手动启动简易HTTP服务器暴露Prometheus指标,这里使用Java 18+自带的SimpleWebServer:
import com.sun.net.httpserver.HttpServer; import java.io.IOException; import java.io.OutputStream; import java.net.InetSocketAddress; public class MetricsServer { public static void start(int port) throws IOException { HttpServer server = HttpServer.create(new InetSocketAddress(port), 0); server.createContext("/metrics", exchange -> { String metrics = MetricsConfig.scrapeMetrics(); exchange.getResponseHeaders().set("Content-Type", "text/plain; version=0.0.4"); exchange.sendResponseHeaders(200, metrics.getBytes().length); try (OutputStream os = exchange.getResponseBody()) { os.write(metrics.getBytes()); } }); server.start(); System.out.println("Metrics server running on port " + port); } public static void main(String[] args) throws IOException { // 启动指标服务器 MetricsServer.start(8080); // 初始化流式处理组件 StreamProcessor processor = new StreamProcessor(); MessageQueue queue = new MessageQueue(); // 模拟消息生产 new Thread(() -> { int i = 0; while (true) { queue.add("message-" + i++); try { Thread.sleep(5); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }).start(); // 模拟消息消费 new Thread(() -> { while (true) { String message = queue.takeMessage(); if (message != null) { processor.processMessage(message); } } }).start(); } }
5. 验证指标
启动应用后,访问http://localhost:8080/metrics,即可看到类似以下的指标输出:
# HELP stream_messages_processed_total Total number of successfully processed messages # TYPE stream_messages_processed_total counter stream_messages_processed_total{application="stream-processing-app",env="production",result="success",} 123.0 # HELP stream_message_processing_duration_seconds Duration of message processing # TYPE stream_message_processing_duration_seconds summary stream_message_processing_duration_seconds_count{application="stream-processing-app",env="production",} 123.0 stream_message_processing_duration_seconds_sum{application="stream-processing-app",env="production",} 1.230
内容的提问来源于stack exchange,提问作者Gautami khante
相关产品推荐
相关产品推荐

