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

如何用Java为Apache Flink编写自定义Prometheus Sink对接Grafana?

Absolutely! You can build a custom Flink Sink to push your processed time-series logs directly to Prometheus, then visualize the data with Grafana. Let’s break this down step by step, including a full Java example tailored to your use case.

Core Approach

Prometheus typically uses a pull-based model, but for streaming data from Flink, the Prometheus PushGateway is the most practical choice—it lets you push metrics from your Flink job to Prometheus, which is perfect for time-series log data. We’ll create a RichSinkFunction (to leverage lifecycle hooks for resource management) that converts processed logs into Prometheus metrics and pushes them to the PushGateway.

Step 1: Add Dependencies

First, include the necessary Prometheus and Flink dependencies in your pom.xml (if using Maven):

<dependencies>
    <!-- Prometheus Java Client -->
    <dependency>
        <groupId>io.prometheus</groupId>
        <artifactId>simpleclient</artifactId>
        <version>0.16.0</version>
    </dependency>
    <!-- Prometheus PushGateway -->
    <dependency>
        <groupId>io.prometheus</groupId>
        <artifactId>simpleclient_pushgateway</artifactId>
        <version>0.16.0</version>
    </dependency>
    <!-- Flink Streaming API -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java</artifactId>
        <version>1.17.0</version>
        <scope>provided</scope>
    </dependency>
</dependencies>

Step 2: Implement the Custom Prometheus Sink

Here’s a complete example of a sink that pushes time-series log metrics (we’ll use a gauge for numeric log fields, but you can adapt this for counters or histograms too):

import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import io.prometheus.client.CollectorRegistry;
import io.prometheus.client.Gauge;
import io.prometheus.client.exporter.PushGateway;

import java.io.IOException;

// Assume your processed time-series log is represented by this POJO
class ProcessedLog {
    private String serviceName;
    private long timestamp;
    private double requestLatency;
    private String requestStatus;

    // Getters and setters
    public String getServiceName() { return serviceName; }
    public void setServiceName(String serviceName) { this.serviceName = serviceName; }
    public long getTimestamp() { return timestamp; }
    public void setTimestamp(long timestamp) { this.timestamp = timestamp; }
    public double getRequestLatency() { return requestLatency; }
    public void setRequestLatency(double requestLatency) { this.requestLatency = requestLatency; }
    public String getRequestStatus() { return requestStatus; }
    public void setRequestStatus(String requestStatus) { this.requestStatus = requestStatus; }
}

public class PrometheusLogSink extends RichSinkFunction<ProcessedLog> {
    private transient PushGateway pushGateway;
    private transient Gauge latencyGauge;
    private final String pushGatewayUrl;
    private final String flinkJobName;

    // Constructor to pass PushGateway configuration
    public PrometheusLogSink(String pushGatewayUrl, String flinkJobName) {
        this.pushGatewayUrl = pushGatewayUrl;
        this.flinkJobName = flinkJobName;
    }

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        // Initialize PushGateway connection
        pushGateway = new PushGateway(pushGatewayUrl);
        // Create a Gauge metric with labels for service name and request status
        latencyGauge = Gauge.build()
                .name("service_request_latency_milliseconds")
                .help("Latency of requests processed by backend services")
                .labelNames("service_name", "request_status")
                .register();
    }

    @Override
    public void invoke(ProcessedLog log, Context context) throws Exception {
        // Set the gauge value using labels from the processed log
        latencyGauge.labels(log.getServiceName(), log.getRequestStatus())
                .set(log.getRequestLatency());
        
        // Push metrics to PushGateway (for high-throughput streams, batch this for better performance)
        try {
            pushGateway.pushAdd(latencyGauge.getCollectorRegistry(), flinkJobName);
        } catch (IOException e) {
            // Track push failures with Flink's built-in metrics
            getRuntimeContext().getMetricGroup().counter("prometheus_push_failures").inc();
            throw new RuntimeException("Failed to push metrics to Prometheus PushGateway", e);
        }
    }

    @Override
    public void close() throws Exception {
        super.close();
        if (pushGateway != null) {
            // Optional: Clean up metrics when the sink shuts down
            pushGateway.delete(flinkJobName, CollectorRegistry.defaultRegistry);
        }
    }
}

Add the custom sink to your Flink streaming pipeline like this:

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class FlinkLogToPrometheusJob {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // Replace this with your actual log processing pipeline
        env.fromElements(
                new ProcessedLog("payment-service", System.currentTimeMillis(), 125.3, "SUCCESS"),
                new ProcessedLog("user-service", System.currentTimeMillis(), 92.1, "SUCCESS"),
                new ProcessedLog("payment-service", System.currentTimeMillis(), 410.7, "ERROR")
        )
        .addSink(new PrometheusLogSink("http://localhost:9091", "flink_time_series_logs")); // PushGateway URL

        env.execute("Flink Time-Series Log Processing to Prometheus");
    }
}

Step 4: Visualize with Grafana

Once your Flink job is pushing metrics to the PushGateway (which Prometheus should be configured to scrape), set up Grafana:

  1. Add Prometheus as a data source (point it to your Prometheus server URL).
  2. Create a new dashboard panel:
    • Select the Prometheus data source.
    • Use a query like service_request_latency_milliseconds{job="flink_time_series_logs"} to pull latency metrics.
    • Choose a visualization type (e.g., Time Series for trend analysis, Gauge for real-time values).

Key Best Practices

  • Batching: For high-throughput streams, avoid pushing metrics on every invoke() call—accumulate metrics in a buffer and push periodically using a timer to reduce overhead.
  • Metric Naming: Follow Prometheus conventions (snake_case, descriptive names) to keep metrics organized.
  • Error Handling: Add retry logic for failed pushes to prevent data loss, and use Flink’s metrics to monitor push success rates.
  • PushGateway Security: Ensure your PushGateway is secured (e.g., with authentication) if running in a production environment.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:09:08