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

Flink如何暴露带运行时动态自定义标签的指标

如何在Flink中实现带运行时动态标签的指标暴露

需求背景

希望在Flink中暴露带运行时动态获取的自定义标签的指标,类似Micrometer的用法:

Metrics.counter("Ops", "type", type).increment(value);

其中type的值仅在运行时确定,直接发布到Prometheus。

现有代码问题

尝试编写了如下RichMapFunction,但由于标签值无法提前预知,无法在open()方法中预先注册指标,每次map()中创建指标会导致重复注册:

public class LogMetrics extends RichMapFunction<IndividualFeatureMetricHolder, String> {

    private MetricGroup metricGroup;

    @Override
    public void open(org.apache.flink.configuration.Configuration parameters) throws Exception {
        super.open(parameters);
        this.metricGroup = getRuntimeContext().getMetricGroup().addGroup("test");
    }

    @Override
    public String map(IndividualFeatureMetricHolder value) throws Exception {
        metricGroup.addGroup("test1", value.getMetricComputationLevel())
                .addGroup("test2", value.getMetricGroup())
                .addGroup("test3", value.getFeature())
                .gauge(value.getMetricType(), value::getValue);
        metricGroup.gauge(value.getMetricType(), value::getValue);
        return "";
    }
}

解决方案

方案1:动态MetricGroup + 指标缓存(推荐)

利用Flink的MetricGroup层级结构动态生成标签,同时缓存已创建的指标避免重复注册:

import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.metrics.Gauge;
import org.apache.flink.metrics.MetricGroup;

import java.util.HashMap;
import java.util.Map;

public class DynamicTagMetrics extends RichMapFunction<IndividualFeatureMetricHolder, String> {

    private MetricGroup rootMetricGroup;
    private Map<String, Gauge<Double>> cachedGauges;

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        // 初始化根指标组
        this.rootMetricGroup = getRuntimeContext().getMetricGroup().addGroup("feature_metrics");
        // 线程安全的缓存,避免多线程下重复创建
        this.cachedGauges = new HashMap<>();
    }

    @Override
    public String map(IndividualFeatureMetricHolder value) throws Exception {
        // 生成唯一标识:组合所有动态标签和指标类型,确保同一标签组合只创建一次指标
        String metricUniqueKey = String.join("|",
                value.getMetricComputationLevel(),
                value.getMetricGroup(),
                value.getFeature(),
                value.getMetricType()
        );

        if (!cachedGauges.containsKey(metricUniqueKey)) {
            // 动态构建带标签的MetricGroup:每个addGroup对应一个Prometheus标签(键值对)
            MetricGroup dynamicTagGroup = rootMetricGroup
                    .addGroup("computation_level", value.getMetricComputationLevel())
                    .addGroup("metric_group", value.getMetricGroup())
                    .addGroup("feature_name", value.getFeature());
            
            // 注册Gauge,用Supplier实时获取最新值
            Gauge<Double> gauge = dynamicTagGroup.gauge(value.getMetricType(), value::getValue);
            cachedGauges.put(metricUniqueKey, gauge);
        }

        // 注:如果指标值是实时更新的,Gauge的Supplier会自动拉取最新值,无需手动更新
        return "";
    }
}

关键说明:

  • Flink的MetricGroup层级会被Prometheus Reporter自动转换为标签,比如上述代码最终生成的指标格式为:feature_metrics_{metricType}{computation_level="xxx", metric_group="xxx", feature_name="xxx"}
  • 缓存避免了重复注册同名指标(Flink不允许同一MetricGroup下重复注册同名指标)
  • 标签基数需控制在合理范围,避免指标爆炸影响监控系统性能

方案2:自定义动态Metric(适用于复杂场景)

如果需要更灵活的多维度指标管理,可以自定义Metric实现,内部维护所有标签组合的数值:

import org.apache.flink.metrics.Gauge;
import java.util.concurrent.ConcurrentHashMap;
import java.util.Map;

// 自定义Gauge,存储不同标签组合的指标值
public class MultiTagGauge implements Gauge<Map<String, Double>> {
    private final ConcurrentHashMap<String, Double> tagValueMap = new ConcurrentHashMap<>();

    public void update(String tagCombination, double value) {
        tagValueMap.put(tagCombination, value);
    }

    @Override
    public Map<String, Double> getValue() {
        return tagValueMap;
    }
}

在Function中使用:

public class CustomDynamicMetrics extends RichMapFunction<IndividualFeatureMetricHolder, String> {

    private MultiTagGauge multiTagGauge;

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        this.multiTagGauge = new MultiTagGauge();
        // 注册自定义Gauge到根指标组
        getRuntimeContext().getMetricGroup().addGroup("custom_metrics")
                .gauge("dynamic_multi_tag", multiTagGauge);
    }

    @Override
    public String map(IndividualFeatureMetricHolder value) throws Exception {
        // 生成标签组合键
        String tagKey = String.format("%s-%s-%s",
                value.getMetricComputationLevel(),
                value.getMetricGroup(),
                value.getFeature()
        );
        // 更新对应标签的指标值
        multiTagGauge.update(tagKey, value.getValue());
        return "";
    }
}

注意:该方案需要对应的MetricReporter支持解析Map类型的Gauge(默认Prometheus Reporter可能不支持,需自定义扩展),因此仅推荐在特殊场景下使用。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 04:20:36