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
相关产品推荐
相关产品推荐

