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

如何在非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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 22:42:24