From 7b2a0b4dac7ba32ca380c945b4110f3c781dfbd3 Mon Sep 17 00:00:00 2001 From: mrproliu <741550557@qq.com> Date: Sat, 17 Oct 2020 22:36:12 +0800 Subject: [PATCH] Provide labeled meter when receive meter (#5659) --- docs/en/setup/backend/backend-meter.md | 5 +- .../meter/config/MeterDataConfig.java | 5 ++ .../meter/process/EvalMultipleData.java | 19 +++++ .../provider/meter/process/MeterBuilder.java | 74 ++++++++++++------- .../meter/process/EvalMultipleDataTest.java | 30 +++++++- .../meter/process/MeterBuilderTest.java | 5 +- .../trace/TraceSampleRateWatcherTest.java | 5 +- 7 files changed, 109 insertions(+), 34 deletions(-) diff --git a/docs/en/setup/backend/backend-meter.md b/docs/en/setup/backend/backend-meter.md index 6d7b85a18..de7f46752 100644 --- a/docs/en/setup/backend/backend-meter.md +++ b/docs/en/setup/backend/backend-meter.md @@ -49,6 +49,9 @@ meters: operation: # Meter value parse groovy script. value: + # Aggregate metrics group by dedicated labels + groupBy: + - # Appoint percentiles if using avgHistogramPercentile operation. percentile: - @@ -56,7 +59,7 @@ meters: #### Meter transform operation -The available operations are `avg`, `avgHistogram` and `avgHistogramPercentile`. The `avg` and `avgXXX` mean to average +The available operations are `avg`, `avgLabeled`, `avgHistogram` and `avgHistogramPercentile`. The `avg` and `avgXXX` mean to average the raw received metrics. When you specify `avgHistogram` and `avgHistogramPercentile`, the source should be the type of `histogram`. diff --git a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/config/MeterDataConfig.java b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/config/MeterDataConfig.java index 828d1348d..c614aeefd 100644 --- a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/config/MeterDataConfig.java +++ b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/config/MeterDataConfig.java @@ -42,4 +42,9 @@ public class MeterDataConfig { */ private List percentile; + /** + * Aggregate meter group by dedicated labels + */ + private List groupBy; + } diff --git a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleData.java b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleData.java index 8ae47a501..a4452333a 100644 --- a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleData.java +++ b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleData.java @@ -22,8 +22,12 @@ import io.vavr.Function2; import java.util.ArrayList; import java.util.List; +import java.util.Map; +import java.util.Objects; import java.util.stream.Collectors; +import static java.util.stream.Collectors.groupingBy; + /** * Combined meter data, has multiple meter data. Support batch process the express on meter */ @@ -122,6 +126,21 @@ public class EvalMultipleData extends EvalData { return dataList.stream().reduce(EvalData::combine).orElseThrow(IllegalArgumentException::new); } + /** + * Combine all of the meter and group by labeled names + * @return Same labeled names and values + */ + Map combineAndGroupBy(List labelNames) { + return dataList.stream() + // group by label + .collect(groupingBy(m -> labelNames.stream().map(l -> m.getLabels().getOrDefault(l, "")).map(Objects::toString).collect(Collectors.joining("-")))) + .entrySet().stream().collect( + // combine labeled values + Collectors.toMap( + e -> e.getKey(), + e -> e.getValue().stream().reduce(EvalData::combine).orElse(null))); + } + /** * Append data to ready eval list */ diff --git a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilder.java b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilder.java index aa0f42149..f8f785361 100644 --- a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilder.java +++ b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilder.java @@ -28,8 +28,12 @@ import org.apache.skywalking.oap.server.core.analysis.meter.MeterSystem; import org.apache.skywalking.oap.server.core.analysis.meter.function.AcceptableValue; import org.apache.skywalking.oap.server.core.analysis.meter.function.AvgHistogramPercentileFunction; import org.apache.skywalking.oap.server.core.analysis.meter.function.BucketedValues; +import org.apache.skywalking.oap.server.core.analysis.metrics.DataTable; +import java.util.Collections; +import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.StringJoiner; import java.util.concurrent.atomic.AtomicBoolean; @@ -42,6 +46,9 @@ public class MeterBuilder { private final MeterConfig config; private final MeterSystem meterSystem; + private final static String DEFAULT_GROUP = "default"; + private final static List DEFAULT_GROUP_LIST = Collections.singletonList(DEFAULT_GROUP); + /** * Current meter has init finished. */ @@ -96,36 +103,49 @@ public class MeterBuilder { log.warn("avg function not support histogram value, please check meter:{}", combinedSingleAvgData.getName()); } break; + case "avgLabeled": + final DataTable dt = new DataTable(); + values.combineAndGroupBy(Optional.ofNullable(config.getMeter().getGroupBy()).orElse(DEFAULT_GROUP_LIST)).entrySet().stream() + .forEach(e -> dt.put(e.getKey(), (long) ((EvalSingleData) e.getValue()).getValue())); + AcceptableValue value = meterSystem.buildMetrics(metricsName, DataTable.class); + value.accept(entity, dt); + value.setTimeBucket(TimeBucket.getMinuteTimeBucket(processor.timestamp())); + meterSystem.doStreamingCalculation(value); + break; case "avgHistogram": case "avgHistogramPercentile": - final EvalData combinedHistogramData = values.combineAsSingleData(); - if (combinedHistogramData instanceof EvalHistogramData) { - final EvalHistogramData histogram = (EvalHistogramData) combinedHistogramData; - long[] buckets = new long[histogram.getBuckets().size()]; - long[] bucketValues = new long[histogram.getBuckets().size()]; - int i = 0; - for (Map.Entry entry : histogram.getBuckets().entrySet()) { - buckets[i] = entry.getKey().intValue(); - bucketValues[i] = entry.getValue(); - i++; - } + values.combineAndGroupBy(Optional.ofNullable(config.getMeter().getGroupBy()).orElse(DEFAULT_GROUP_LIST)).entrySet().stream() + .forEach(e -> { + final String group = e.getKey(); + final EvalData combinedHistogramData = e.getValue(); + if (combinedHistogramData instanceof EvalHistogramData) { + final EvalHistogramData histogram = (EvalHistogramData) combinedHistogramData; + long[] buckets = new long[histogram.getBuckets().size()]; + long[] bucketValues = new long[histogram.getBuckets().size()]; + int i = 0; + for (Map.Entry entry : histogram.getBuckets().entrySet()) { + buckets[i] = entry.getKey().intValue(); + bucketValues[i] = entry.getValue(); + i++; + } - if (config.getMeter().getOperation().equals("avgHistogram")) { - AcceptableValue avgHistogramValue = meterSystem.buildMetrics(metricsName, BucketedValues.class); - avgHistogramValue.accept(entity, new BucketedValues(buckets, bucketValues)); - avgHistogramValue.setTimeBucket(TimeBucket.getMinuteTimeBucket(processor.timestamp())); - meterSystem.doStreamingCalculation(avgHistogramValue); - } else { - final AcceptableValue percentileValue = - meterSystem.buildMetrics(metricsName, AvgHistogramPercentileFunction.AvgPercentileArgument.class); - percentileValue.accept(entity, new AvgHistogramPercentileFunction.AvgPercentileArgument(new BucketedValues(buckets, bucketValues), - config.getMeter().getPercentile().stream().mapToInt(Integer::intValue).toArray())); - percentileValue.setTimeBucket(TimeBucket.getMinuteTimeBucket(processor.timestamp())); - meterSystem.doStreamingCalculation(percentileValue); - } - } else { - log.warn(config.getMeter().getOperation() + " function not support single value, please check meter:{}", combinedHistogramData.getName()); - } + final BucketedValues bucketedValues = new BucketedValues(buckets, bucketValues); + bucketedValues.setGroup(group); + if (config.getMeter().getOperation().equals("avgHistogram")) { + AcceptableValue avgHistogramValue = meterSystem.buildMetrics(metricsName, BucketedValues.class); + avgHistogramValue.accept(entity, bucketedValues); + avgHistogramValue.setTimeBucket(TimeBucket.getMinuteTimeBucket(processor.timestamp())); + meterSystem.doStreamingCalculation(avgHistogramValue); + } else { + final AcceptableValue percentileValue = + meterSystem.buildMetrics(metricsName, AvgHistogramPercentileFunction.AvgPercentileArgument.class); + percentileValue.accept(entity, new AvgHistogramPercentileFunction.AvgPercentileArgument(bucketedValues, + config.getMeter().getPercentile().stream().mapToInt(Integer::intValue).toArray())); + percentileValue.setTimeBucket(TimeBucket.getMinuteTimeBucket(processor.timestamp())); + meterSystem.doStreamingCalculation(percentileValue); + } + } + }); break; default: log.warn("Cannot support function:{}", config.getMeter().getOperation()); diff --git a/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleDataTest.java b/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleDataTest.java index 5e1965ee6..dbbb2186a 100644 --- a/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleDataTest.java +++ b/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/EvalMultipleDataTest.java @@ -19,13 +19,15 @@ package org.apache.skywalking.oap.server.analyzer.provider.meter.process; import io.vavr.Function2; -import java.util.Arrays; -import java.util.List; import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.powermock.reflect.Whitebox; +import java.util.Arrays; +import java.util.List; +import java.util.Map; + public class EvalMultipleDataTest extends EvalDataBaseTest { private EvalMultipleData singleMultiple; @@ -104,6 +106,30 @@ public class EvalMultipleDataTest extends EvalDataBaseTest { Assert.assertEquals(35, combinedHistogramData.getBuckets().get(10d).longValue()); } + @Test + public void testCombineAndGroupBy() { + final Map singleValueCombineAndGroupBy = singleMultiple.combineAndGroupBy(Arrays.asList("k1")); + Assert.assertEquals(3, singleValueCombineAndGroupBy.size()); + Assert.assertEquals(10d, ((EvalSingleData) singleValueCombineAndGroupBy.get("v1")).getValue(), 0.0); + Assert.assertEquals(20d, ((EvalSingleData) singleValueCombineAndGroupBy.get("v2")).getValue(), 0.0); + Assert.assertEquals(30d, ((EvalSingleData) singleValueCombineAndGroupBy.get("v3")).getValue(), 0.0); + + final Map histogramCombineAndGroupBy = histogramMultiple.combineAndGroupBy(Arrays.asList("k1")); + Assert.assertEquals(3, histogramCombineAndGroupBy.size()); + Assert.assertEquals(10, ((EvalHistogramData) histogramCombineAndGroupBy.get("v1")).getBuckets().get(1d).longValue()); + Assert.assertEquals(20, ((EvalHistogramData) histogramCombineAndGroupBy.get("v1")).getBuckets().get(5d).longValue()); + Assert.assertEquals(3, ((EvalHistogramData) histogramCombineAndGroupBy.get("v1")).getBuckets().get(10d).longValue()); + + Assert.assertEquals(5, ((EvalHistogramData) histogramCombineAndGroupBy.get("v2")).getBuckets().get(1d).longValue()); + Assert.assertEquals(10, ((EvalHistogramData) histogramCombineAndGroupBy.get("v2")).getBuckets().get(5d).longValue()); + Assert.assertEquals(7, ((EvalHistogramData) histogramCombineAndGroupBy.get("v2")).getBuckets().get(10d).longValue()); + + Assert.assertEquals(15, ((EvalHistogramData) histogramCombineAndGroupBy.get("v3")).getBuckets().get(1d).longValue()); + Assert.assertEquals(20, ((EvalHistogramData) histogramCombineAndGroupBy.get("v3")).getBuckets().get(5d).longValue()); + Assert.assertEquals(25, ((EvalHistogramData) histogramCombineAndGroupBy.get("v3")).getBuckets().get(10d).longValue()); + + } + /** * Verify the multiple data operation */ diff --git a/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilderTest.java b/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilderTest.java index 8a4a63efa..b4d4455d0 100644 --- a/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilderTest.java +++ b/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/meter/process/MeterBuilderTest.java @@ -18,8 +18,6 @@ package org.apache.skywalking.oap.server.analyzer.provider.meter.process; -import java.util.ArrayList; -import java.util.List; import org.apache.skywalking.oap.server.core.analysis.IDManager; import org.apache.skywalking.oap.server.core.analysis.TimeBucket; import org.apache.skywalking.oap.server.core.analysis.meter.function.AcceptableValue; @@ -35,6 +33,9 @@ import org.powermock.core.classloader.annotations.PowerMockIgnore; import org.powermock.modules.junit4.PowerMockRunner; import org.powermock.reflect.Whitebox; +import java.util.ArrayList; +import java.util.List; + import static org.mockito.Matchers.any; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doCallRealMethod; diff --git a/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/trace/TraceSampleRateWatcherTest.java b/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/trace/TraceSampleRateWatcherTest.java index a1584b373..2a16e7801 100644 --- a/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/trace/TraceSampleRateWatcherTest.java +++ b/oap-server/analyzer/agent-analyzer/src/test/java/org/apache/skywalking/oap/server/analyzer/provider/trace/TraceSampleRateWatcherTest.java @@ -18,8 +18,6 @@ package org.apache.skywalking.oap.server.analyzer.provider.trace; -import java.util.Optional; -import java.util.Set; import org.apache.skywalking.oap.server.analyzer.provider.AnalyzerModuleProvider; import org.apache.skywalking.oap.server.configuration.api.ConfigChangeWatcher; import org.apache.skywalking.oap.server.configuration.api.ConfigTable; @@ -30,6 +28,9 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.runners.MockitoJUnitRunner; +import java.util.Optional; +import java.util.Set; + import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.MatcherAssert.assertThat;