This is an automated email from the ASF dual-hosted git repository.
wusheng pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/skywalking.git
The following commit(s) were added to refs/heads/master by this push:
new 6a78773743 Improve BanyanDB measure schema (#10233)
6a78773743 is described below
commit 6a7877374395f7a99365873bd1e033c77019bf12
Author: Gao Hongtao <[email protected]>
AuthorDate: Sun Jan 8 18:50:00 2023 +0800
Improve BanyanDB measure schema (#10233)
---
docs/en/changes/changes.md | 6 +-
oap-server-bom/pom.xml | 2 +-
.../analysis/manual/endpoint/EndpointTraffic.java | 3 +
.../analysis/manual/instance/InstanceTraffic.java | 3 +
.../analysis/manual/process/ProcessTraffic.java | 4 +
.../ServiceInstanceRelationClientSideMetrics.java | 2 +
.../manual/searchtag/TagAutocompleteData.java | 4 +
.../analysis/manual/service/ServiceTraffic.java | 3 +
.../analysis/meter/function/HistogramFunction.java | 1 +
.../meter/function/PercentileFunction.java | 3 +
.../analysis/meter/function/avg/AvgFunction.java | 3 +
.../meter/function/avg/AvgHistogramFunction.java | 3 +
.../avg/AvgHistogramPercentileFunction.java | 4 +
.../meter/function/avg/AvgLabeledFunction.java | 3 +
.../meter/function/latest/LatestFunction.java | 1 +
.../analysis/meter/function/sum/SumFunction.java | 1 +
.../sum/SumHistogramPercentileFunction.java | 3 +
.../function/sumpermin/SumPerMinFunction.java | 2 +
.../sumpermin/SumPerMinLabeledFunction.java | 2 +
.../server/core/analysis/metrics/ApdexMetrics.java | 5 +
.../server/core/analysis/metrics/CPMMetrics.java | 3 +
.../server/core/analysis/metrics/CountMetrics.java | 2 +
.../core/analysis/metrics/DoubleAvgMetrics.java | 4 +
.../oap/server/core/analysis/metrics/Event.java | 2 +
.../core/analysis/metrics/HistogramMetrics.java | 2 +
.../core/analysis/metrics/LongAvgMetrics.java | 4 +
.../core/analysis/metrics/MaxDoubleMetrics.java | 2 +
.../core/analysis/metrics/MaxLongMetrics.java | 2 +
.../core/analysis/metrics/MinDoubleMetrics.java | 2 +
.../core/analysis/metrics/MinLongMetrics.java | 2 +
.../core/analysis/metrics/PercentMetrics.java | 5 +-
.../core/analysis/metrics/PercentileMetrics.java | 4 +
.../server/core/analysis/metrics/RateMetrics.java | 4 +
.../server/core/analysis/metrics/SumMetrics.java | 2 +
.../ebpf/storage/EBPFProfilingScheduleRecord.java | 2 +
.../oap/server/core/source/ScopeDefaultColumn.java | 3 +-
.../server/core/storage/annotation/BanyanDB.java | 33 ++++++
.../core/storage/model/BanyanDBExtension.java | 6 +
.../core/storage/model/BanyanDBModelExtension.java | 8 ++
.../server/core/storage/model/StorageModels.java | 13 ++-
.../core/zipkin/ZipkinServiceRelationTraffic.java | 3 +
.../core/zipkin/ZipkinServiceSpanTraffic.java | 3 +
.../server/core/zipkin/ZipkinServiceTraffic.java | 2 +
.../oap/server/core/zipkin/ZipkinSpanRecord.java | 3 +-
.../server/core/storage/model/ModelColumnTest.java | 10 +-
.../storage/plugin/banyandb/BanyanDBConverter.java | 5 +-
.../plugin/banyandb/BanyanDBIndexInstaller.java | 2 +-
.../storage/plugin/banyandb/MetadataRegistry.java | 103 +++++++++--------
.../BanyanDBEBPFProfilingScheduleQueryDAO.java | 67 +++++------
.../banyandb/measure/BanyanDBMetadataQueryDAO.java | 120 +++++++++++---------
.../banyandb/measure/BanyanDBMetricsDAO.java | 123 +++++++++++++++++++--
.../banyandb/measure/BanyanDBMetricsQueryDAO.java | 92 ++++++++-------
.../banyandb/stream/AbstractBanyanDBDAO.java | 5 -
.../stream/BanyanDBProfileTaskLogQueryDAO.java | 8 +-
.../BanyanDBProfileThreadSnapshotQueryDAO.java | 3 +-
.../banyandb/stream/BanyanDBTraceQueryDAO.java | 3 +-
skywalking-ui | 2 +-
.../java-test-service/e2e-protocol/src/main/proto | 2 +-
test/e2e-v2/script/env | 2 +-
59 files changed, 503 insertions(+), 218 deletions(-)
diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md
index 36a37b9ce7..2d8eb04bd2 100644
--- a/docs/en/changes/changes.md
+++ b/docs/en/changes/changes.md
@@ -69,7 +69,10 @@
* Correct the TopN record query DAO of BanyanDB.
* Tweak interval settings of BanyanDB.
* Support monitoring AWS Cloud EKS.
-* Bump BanyanDB Java client to 0.3.0-rc0.
+* Bump BanyanDB Java client to 0.3.0-rc1.
+* Remove `id` tag from measures.
+* Add `Banyandb.MeasureField` to mark a column as a BanyanDB Measure field.
+* Add `BanyanDB.StoreIDTag` to store a process's id for searching.
* [**Breaking Change**] The supported version of ShardingSphere-Proxy is
upgraded from 5.1.2 to 5.3.1. Due to the changes of ShardingSphere's API,
versions before 5.3.1 are not compatible.
* Add the eBPF network profiling E2E Test in the per storage.
@@ -81,7 +84,6 @@
* Add Micrometer icon
* Update MySQL UI to support MariaDB
* Add AWS menu for supporting AWS monitoring
-* Add Zipkin Lens UI as part of Service Mesh dashboard.
#### Documentation
diff --git a/oap-server-bom/pom.xml b/oap-server-bom/pom.xml
index 9d3a88baa1..3d0f63467d 100644
--- a/oap-server-bom/pom.xml
+++ b/oap-server-bom/pom.xml
@@ -73,7 +73,7 @@
<awaitility.version>3.0.0</awaitility.version>
<httpcore.version>4.4.13</httpcore.version>
<commons-compress.version>1.21</commons-compress.version>
- <banyandb-java-client.version>0.3.0-rc0</banyandb-java-client.version>
+ <banyandb-java-client.version>0.3.0-rc1</banyandb-java-client.version>
<kafka-clients.version>2.8.1</kafka-clients.version>
<spring-kafka-test.version>2.4.6.RELEASE</spring-kafka-test.version>
</properties>
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpoint/EndpointTraffic.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpoint/EndpointTraffic.java
index a520cb7d14..3c764e7b9a 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpoint/EndpointTraffic.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/endpoint/EndpointTraffic.java
@@ -32,6 +32,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
@@ -57,12 +58,14 @@ public class EndpointTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = SERVICE_ID)
+ @BanyanDB.SeriesID(index = 0)
private String serviceId;
@Setter
@Getter
@Column(columnName = NAME)
@ElasticSearch.MatchQuery
@ElasticSearch.Column(columnAlias = "endpoint_traffic_name")
+ @BanyanDB.SeriesID(index = 1)
private String name = Const.EMPTY_STRING;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/instance/InstanceTraffic.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/instance/InstanceTraffic.java
index d69af09e5e..eb3e1bc45a 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/instance/InstanceTraffic.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/instance/InstanceTraffic.java
@@ -32,6 +32,7 @@ import
org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
@@ -62,12 +63,14 @@ public class InstanceTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = SERVICE_ID)
+ @BanyanDB.SeriesID(index = 0)
private String serviceId;
@Setter
@Getter
@Column(columnName = NAME, storageOnly = true)
@ElasticSearch.Column(columnAlias = "instance_traffic_name")
+ @BanyanDB.SeriesID(index = 1)
private String name;
@Setter
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/process/ProcessTraffic.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/process/ProcessTraffic.java
index 98e14a4336..174d873702 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/process/ProcessTraffic.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/process/ProcessTraffic.java
@@ -34,6 +34,7 @@ import
org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -51,6 +52,7 @@ import static
org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.PR
"name",
})
@SQLDatabase.Sharding(shardingAlgorithm = ShardingAlgorithm.NO_SHARDING)
[email protected]
public class ProcessTraffic extends Metrics {
public static final String INDEX_NAME = "process_traffic";
public static final String SERVICE_ID = "service_id";
@@ -73,6 +75,7 @@ public class ProcessTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = INSTANCE_ID, length = 600)
+ @BanyanDB.SeriesID(index = 0)
private String instanceId;
@Getter
@@ -82,6 +85,7 @@ public class ProcessTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = NAME, length = 500)
+ @BanyanDB.SeriesID(index = 1)
private String name;
@Setter
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/relation/instance/ServiceInstanceRelationClientSideMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/relation/instance/ServiceInstanceRelationClientSideMetrics.java
index 17c60b848b..5304666d51 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/relation/instance/ServiceInstanceRelationClientSideMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/relation/instance/ServiceInstanceRelationClientSideMetrics.java
@@ -29,6 +29,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -72,6 +73,7 @@ public class ServiceInstanceRelationClientSideMetrics extends
Metrics {
@Setter
@Getter
@Column(columnName = ENTITY_ID, length = 512)
+ @BanyanDB.SeriesID(index = 0)
private String entityId;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/searchtag/TagAutocompleteData.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/searchtag/TagAutocompleteData.java
index 7474fd24d8..cfc10c9647 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/searchtag/TagAutocompleteData.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/searchtag/TagAutocompleteData.java
@@ -29,6 +29,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -56,15 +57,18 @@ public class TagAutocompleteData extends Metrics {
@Setter
@Getter
@Column(columnName = TAG_KEY)
+ @BanyanDB.SeriesID(index = 1)
private String tagKey;
@Setter
@Getter
@Column(columnName = TAG_VALUE, length = Tag.TAG_LENGTH)
+ @BanyanDB.SeriesID(index = 2)
private String tagValue;
@Setter
@Getter
@Column(columnName = TAG_TYPE)
+ @BanyanDB.SeriesID(index = 0)
private String tagType;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/service/ServiceTraffic.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/service/ServiceTraffic.java
index 3aaca1e99d..cb303244a8 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/service/ServiceTraffic.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/service/ServiceTraffic.java
@@ -32,6 +32,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
@@ -68,6 +69,7 @@ public class ServiceTraffic extends Metrics {
@Column(columnName = NAME)
@ElasticSearch.MatchQuery
@ElasticSearch.Column(columnAlias = "service_traffic_name")
+ @BanyanDB.SeriesID(index = 1)
private String name = Const.EMPTY_STRING;
@Setter
@@ -90,6 +92,7 @@ public class ServiceTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = LAYER)
+ @BanyanDB.SeriesID(index = 0)
private Layer layer = Layer.UNDEFINED;
/**
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/HistogramFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/HistogramFunction.java
index 6f01b1e628..cb90a6cb40 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/HistogramFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/HistogramFunction.java
@@ -55,6 +55,7 @@ public abstract class HistogramFunction extends Meter
implements AcceptableValue
@Getter
@Setter
@Column(columnName = DATASET, dataType = Column.ValueDataType.HISTOGRAM,
storageOnly = true, defaultValue = 0)
+ @BanyanDB.MeasureField
private DataTable dataset = new DataTable(30);
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/PercentileFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/PercentileFunction.java
index c8961eac1b..3279c55791 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/PercentileFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/PercentileFunction.java
@@ -64,10 +64,12 @@ public abstract class PercentileFunction extends Meter
implements AcceptableValu
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.LABELED_VALUE,
storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_value")
+ @BanyanDB.MeasureField
private DataTable percentileValues = new DataTable(10);
@Getter
@Setter
@Column(columnName = DATASET, storageOnly = true)
+ @BanyanDB.MeasureField
private DataTable dataset = new DataTable(30);
/**
* Rank
@@ -75,6 +77,7 @@ public abstract class PercentileFunction extends Meter
implements AcceptableValu
@Getter
@Setter
@Column(columnName = RANKS, storageOnly = true)
+ @BanyanDB.MeasureField
private IntList ranks = new IntList(10);
private boolean isCalculated = false;
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgFunction.java
index f01ff10ac2..d7aed9e3c4 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgFunction.java
@@ -66,14 +66,17 @@ public abstract class AvgFunction extends Meter implements
AcceptableValue<Long>
@Getter
@Setter
@Column(columnName = SUMMATION, storageOnly = true)
+ @BanyanDB.MeasureField
protected long summation;
@Getter
@Setter
@Column(columnName = COUNT, storageOnly = true)
+ @BanyanDB.MeasureField
protected long count;
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Avg)
+ @BanyanDB.MeasureField
private long value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramFunction.java
index 46289e676d..96a6bc47d2 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramFunction.java
@@ -68,15 +68,18 @@ public abstract class AvgHistogramFunction extends Meter
implements AcceptableVa
@Setter
@Column(columnName = SUMMATION, storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_summation")
+ @BanyanDB.MeasureField
protected DataTable summation = new DataTable(30);
@Getter
@Setter
@Column(columnName = COUNT, storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_count")
+ @BanyanDB.MeasureField
protected DataTable count = new DataTable(30);
@Getter
@Setter
@Column(columnName = DATASET, dataType = Column.ValueDataType.HISTOGRAM,
storageOnly = true, defaultValue = 0)
+ @BanyanDB.MeasureField
private DataTable dataset = new DataTable(30);
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramPercentileFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramPercentileFunction.java
index 186e047b02..1c3d12068d 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramPercentileFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgHistogramPercentileFunction.java
@@ -83,20 +83,24 @@ public abstract class AvgHistogramPercentileFunction
extends Meter implements Ac
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.LABELED_VALUE,
storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_value")
+ @BanyanDB.MeasureField
private DataTable percentileValues = new DataTable(10);
@Getter
@Setter
@Column(columnName = SUMMATION, storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_summation")
+ @BanyanDB.MeasureField
protected DataTable summation = new DataTable(30);
@Getter
@Setter
@Column(columnName = COUNT, storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_count")
+ @BanyanDB.MeasureField
protected DataTable count = new DataTable(30);
@Getter
@Setter
@Column(columnName = DATASET, storageOnly = true)
+ @BanyanDB.MeasureField
private DataTable dataset = new DataTable(30);
/**
* Rank
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgLabeledFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgLabeledFunction.java
index 9e78a41da8..c7c46f4e6f 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgLabeledFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/avg/AvgLabeledFunction.java
@@ -66,16 +66,19 @@ public abstract class AvgLabeledFunction extends Meter
implements AcceptableValu
@Setter
@Column(columnName = SUMMATION, storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_summation")
+ @BanyanDB.MeasureField
protected DataTable summation = new DataTable(30);
@Getter
@Setter
@Column(columnName = COUNT, storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_count")
+ @BanyanDB.MeasureField
protected DataTable count = new DataTable(30);
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.LABELED_VALUE,
storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_value")
+ @BanyanDB.MeasureField
private DataTable value = new DataTable(30);
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/latest/LatestFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/latest/LatestFunction.java
index 800320048c..c5145212b6 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/latest/LatestFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/latest/LatestFunction.java
@@ -63,6 +63,7 @@ public abstract class LatestFunction extends Meter implements
AcceptableValue<Lo
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Latest)
+ @BanyanDB.MeasureField
private long value;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumFunction.java
index 43c3012ab1..e973b07313 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumFunction.java
@@ -60,6 +60,7 @@ public abstract class SumFunction extends Meter implements
AcceptableValue<Long>
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Sum)
+ @BanyanDB.MeasureField
private long value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumHistogramPercentileFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumHistogramPercentileFunction.java
index 305a7e3d70..11d8379603 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumHistogramPercentileFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sum/SumHistogramPercentileFunction.java
@@ -73,11 +73,13 @@ public abstract class SumHistogramPercentileFunction
extends Meter implements Ac
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.LABELED_VALUE,
storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_value")
+ @BanyanDB.MeasureField
private DataTable percentileValues = new DataTable(10);
@Getter
@Setter
@Column(columnName = SUMMATION, storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_summation")
+ @BanyanDB.MeasureField
protected DataTable summation = new DataTable(30);
/**
* Rank
@@ -85,6 +87,7 @@ public abstract class SumHistogramPercentileFunction extends
Meter implements Ac
@Getter
@Setter
@Column(columnName = RANKS, storageOnly = true)
+ @BanyanDB.MeasureField
private IntList ranks = new IntList(10);
private boolean isCalculated = false;
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinFunction.java
index 3e53c18c73..77b439476b 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinFunction.java
@@ -61,11 +61,13 @@ public abstract class SumPerMinFunction extends Meter
implements AcceptableValue
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Avg)
+ @BanyanDB.MeasureField
private long value;
@Getter
@Setter
@Column(columnName = TOTAL, storageOnly = true)
+ @BanyanDB.MeasureField
private long total;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinLabeledFunction.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinLabeledFunction.java
index 50919f43d3..7acfd06c66 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinLabeledFunction.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/function/sumpermin/SumPerMinLabeledFunction.java
@@ -60,11 +60,13 @@ public abstract class SumPerMinLabeledFunction extends
Meter implements Acceptab
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.LABELED_VALUE,
storageOnly = true)
+ @BanyanDB.MeasureField
private DataTable value = new DataTable(30);
@Getter
@Setter
@Column(columnName = TOTAL, storageOnly = true)
+ @BanyanDB.MeasureField
private DataTable total = new DataTable(30);
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/ApdexMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/ApdexMetrics.java
index 873a153def..a821d6574b 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/ApdexMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/ApdexMetrics.java
@@ -26,6 +26,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
@@ -50,18 +51,22 @@ public abstract class ApdexMetrics extends Metrics
implements IntValueHolder {
@Getter
@Setter
@Column(columnName = TOTAL_NUM, storageOnly = true)
+ @BanyanDB.MeasureField
private long totalNum;
@Getter
@Setter
@Column(columnName = S_NUM, storageOnly = true)
+ @BanyanDB.MeasureField
private long sNum;
@Getter
@Setter
@Column(columnName = T_NUM, storageOnly = true)
+ @BanyanDB.MeasureField
private long tNum;
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Avg)
+ @BanyanDB.MeasureField
@ElasticSearch.Column(columnAlias = "int_value")
private int value;
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CPMMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CPMMetrics.java
index bc0f02f01f..5428a344eb 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CPMMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CPMMetrics.java
@@ -24,6 +24,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.ConstOn
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "cpm")
@@ -35,10 +36,12 @@ public abstract class CPMMetrics extends Metrics implements
LongValueHolder {
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Avg)
+ @BanyanDB.MeasureField
private long value;
@Getter
@Setter
@Column(columnName = TOTAL, storageOnly = true)
+ @BanyanDB.MeasureField
private long total;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CountMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CountMetrics.java
index 9a223147f3..a044837717 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CountMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/CountMetrics.java
@@ -24,6 +24,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.ConstOn
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "count")
@@ -34,6 +35,7 @@ public abstract class CountMetrics extends Metrics implements
LongValueHolder {
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Sum)
+ @BanyanDB.MeasureField
private long value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/DoubleAvgMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/DoubleAvgMetrics.java
index 1e9d7b742a..bce29c98bf 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/DoubleAvgMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/DoubleAvgMetrics.java
@@ -25,6 +25,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
@@ -39,15 +40,18 @@ public abstract class DoubleAvgMetrics extends Metrics
implements DoubleValueHol
@Setter
@Column(columnName = SUMMATION, storageOnly = true)
@ElasticSearch.Column(columnAlias = "double_summation")
+ @BanyanDB.MeasureField
private double summation;
@Getter
@Setter
@Column(columnName = COUNT, storageOnly = true)
+ @BanyanDB.MeasureField
private long count;
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Avg)
@ElasticSearch.Column(columnAlias = "double_value")
+ @BanyanDB.MeasureField
private double value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/Event.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/Event.java
index cc14d37545..12a434812f 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/Event.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/Event.java
@@ -31,6 +31,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.ScopeDeclaration;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -84,6 +85,7 @@ public class Event extends Metrics {
}
@Column(columnName = UUID)
+ @BanyanDB.SeriesID(index = 0)
private String uuid;
@Column(columnName = SERVICE)
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/HistogramMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/HistogramMetrics.java
index dbb4870857..9e6fa6f01c 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/HistogramMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/HistogramMetrics.java
@@ -24,6 +24,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Arg;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
/**
@@ -42,6 +43,7 @@ public abstract class HistogramMetrics extends Metrics {
@Getter
@Setter
@Column(columnName = DATASET, dataType = Column.ValueDataType.HISTOGRAM,
storageOnly = true, defaultValue = 0)
+ @BanyanDB.MeasureField
private DataTable dataset = new DataTable(30);
/**
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/LongAvgMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/LongAvgMetrics.java
index 661281c770..dd9bb25d95 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/LongAvgMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/LongAvgMetrics.java
@@ -25,6 +25,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "longAvg")
@@ -37,14 +38,17 @@ public abstract class LongAvgMetrics extends Metrics
implements LongValueHolder
@Getter
@Setter
@Column(columnName = SUMMATION, storageOnly = true)
+ @BanyanDB.MeasureField
protected long summation;
@Getter
@Setter
@Column(columnName = COUNT, storageOnly = true)
+ @BanyanDB.MeasureField
protected long count;
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Avg)
+ @BanyanDB.MeasureField
private long value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxDoubleMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxDoubleMetrics.java
index 5dcef2369c..8929e67dfd 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxDoubleMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxDoubleMetrics.java
@@ -23,6 +23,7 @@ import lombok.Setter;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "maxDouble")
@@ -33,6 +34,7 @@ public abstract class MaxDoubleMetrics extends Metrics
implements DoubleValueHol
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE)
+ @BanyanDB.MeasureField
private double value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxLongMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxLongMetrics.java
index 0d3043faca..5785e72d51 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxLongMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MaxLongMetrics.java
@@ -23,6 +23,7 @@ import lombok.Setter;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "max")
@@ -33,6 +34,7 @@ public abstract class MaxLongMetrics extends Metrics
implements LongValueHolder
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE)
+ @BanyanDB.MeasureField
private long value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinDoubleMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinDoubleMetrics.java
index bdc45e63e8..005e972f8d 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinDoubleMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinDoubleMetrics.java
@@ -23,6 +23,7 @@ import lombok.Setter;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "minDouble")
@@ -33,6 +34,7 @@ public abstract class MinDoubleMetrics extends Metrics
implements DoubleValueHol
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE)
+ @BanyanDB.MeasureField
private double value = Double.MAX_VALUE;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinLongMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinLongMetrics.java
index 88f7fab2c5..61b8ef587f 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinLongMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/MinLongMetrics.java
@@ -23,6 +23,7 @@ import lombok.Setter;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "min")
@@ -33,6 +34,7 @@ public abstract class MinLongMetrics extends Metrics
implements LongValueHolder
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE)
+ @BanyanDB.MeasureField
private long value = Long.MAX_VALUE;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentMetrics.java
index 2efd2a6ee3..0d8c0475be 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentMetrics.java
@@ -24,6 +24,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Expression;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "percent")
@@ -35,14 +36,16 @@ public abstract class PercentMetrics extends Metrics
implements IntValueHolder {
@Getter
@Setter
@Column(columnName = TOTAL, storageOnly = true)
+ @BanyanDB.MeasureField
private long total;
@Getter
@Setter
@Column(columnName = PERCENTAGE, dataType =
Column.ValueDataType.COMMON_VALUE, function = Function.Avg)
+ @BanyanDB.MeasureField
private int percentage;
@Getter
@Setter
- @Column(columnName = MATCH)
+ @Column(columnName = MATCH, storageOnly = true)
private long match;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentileMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentileMetrics.java
index 24aa7b71f5..b58d490fa5 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentileMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/PercentileMetrics.java
@@ -27,6 +27,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Arg;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entrance;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearch;
@@ -54,14 +55,17 @@ public abstract class PercentileMetrics extends Metrics
implements MultiIntValue
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.LABELED_VALUE,
storageOnly = true)
@ElasticSearch.Column(columnAlias = "datatable_value")
+ @BanyanDB.MeasureField
private DataTable percentileValues;
@Getter
@Setter
@Column(columnName = PRECISION, storageOnly = true)
+ @BanyanDB.MeasureField
private int precision;
@Getter
@Setter
@Column(columnName = DATASET, storageOnly = true)
+ @BanyanDB.MeasureField
private DataTable dataset;
private boolean isCalculated;
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/RateMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/RateMetrics.java
index 830f1d7c5e..f107a78f3c 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/RateMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/RateMetrics.java
@@ -23,6 +23,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Expression;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "rate")
@@ -34,14 +35,17 @@ public abstract class RateMetrics extends Metrics
implements IntValueHolder {
@Getter
@Setter
@Column(columnName = DENOMINATOR)
+ @BanyanDB.MeasureField
private long denominator;
@Getter
@Setter
@Column(columnName = PERCENTAGE, dataType =
Column.ValueDataType.COMMON_VALUE, function = Function.Avg)
+ @BanyanDB.MeasureField
private int percentage;
@Getter
@Setter
@Column(columnName = NUMERATOR)
+ @BanyanDB.MeasureField
private long numerator;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/SumMetrics.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/SumMetrics.java
index a2f3d61fb9..7b1cb47bee 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/SumMetrics.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/metrics/SumMetrics.java
@@ -24,6 +24,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.Entranc
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.MetricsFunction;
import
org.apache.skywalking.oap.server.core.analysis.metrics.annotation.SourceFrom;
import org.apache.skywalking.oap.server.core.query.sql.Function;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
@MetricsFunction(functionName = "sum")
@@ -34,6 +35,7 @@ public abstract class SumMetrics extends Metrics implements
LongValueHolder {
@Getter
@Setter
@Column(columnName = VALUE, dataType = Column.ValueDataType.COMMON_VALUE,
function = Function.Sum)
+ @BanyanDB.MeasureField
private long value;
@Entrance
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingScheduleRecord.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingScheduleRecord.java
index 2c086d488e..6af14c30ff 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingScheduleRecord.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingScheduleRecord.java
@@ -28,6 +28,7 @@ import
org.apache.skywalking.oap.server.core.analysis.worker.MetricsStreamProces
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -70,6 +71,7 @@ public class EBPFProfilingScheduleRecord extends Metrics {
@Column(columnName = END_TIME)
private long endTime;
@Column(columnName = EBPF_PROFILING_SCHEDULE_ID)
+ @BanyanDB.SeriesID(index = 0)
private String scheduleId;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/ScopeDefaultColumn.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/ScopeDefaultColumn.java
index c24706c98f..61566f9f70 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/ScopeDefaultColumn.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/ScopeDefaultColumn.java
@@ -52,7 +52,8 @@ public class ScopeDefaultColumn {
/**
* Dynamic active means this column is only activated through core
setting explicitly.
*
- * @return
+ * @return FALSE: this column is not going to be added to the final
generated metric as a column.
+ * TRUE: this column could be added as a column if
core/activeExtraModelColumns == true.
*/
boolean requireDynamicActive() default false;
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/BanyanDB.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/BanyanDB.java
index d3cd601081..9777c93263 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/BanyanDB.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/BanyanDB.java
@@ -76,6 +76,14 @@ public @interface BanyanDB {
/**
* Relative entity tag
*
+ * The index number determines the order of the column placed in the
SeriesID.
+ * BanyanDB SeriesID searching procedure uses a prefix-scanning
strategy.
+ * Searching series against a prefix could improve the performance.
+ * <p>
+ * For example, the ServiceTraffic composite "layer" and "name" as the
SeriesID,
+ * considering OAP finds services by "layer", the "layer" 's index
should be 0 to
+ * trigger a prefix-scanning.
+ *
* @return index, from zero.
*/
int index() default -1;
@@ -131,4 +139,29 @@ public @interface BanyanDB {
@interface TimestampColumn {
String value();
}
+
+ /**
+ * MeasureField defines a column as a measure's field.
+ *
+ * Annotated: the column is a measure field.
+ * Unannotated: the column is a measure tag.
+ * storageOnly=true: the column is a measure tag that is not indexed.
+ * storageOnly=false: the column is a measure tag that is indexed.
+ * indexOnly=true: the column is a measure tag that is indexed, but not
stored.
+ * indexOnly=false: the column is a measure tag that is indexed and
stored.
+ * @since 9.4.0
+ */
+ @Target({ElementType.FIELD})
+ @Retention(RetentionPolicy.RUNTIME)
+ @interface MeasureField {
+ }
+
+ /**
+ * StoreIDTag indicates a metric store its ID as a tag for searching.
+ * @Since 9.4.0
+ */
+ @Target({ElementType.TYPE})
+ @Retention(RetentionPolicy.RUNTIME)
+ @interface StoreIDAsTag {
+ }
}
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBExtension.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBExtension.java
index 5c2ab25438..0c0ed11cdb 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBExtension.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBExtension.java
@@ -63,6 +63,12 @@ public class BanyanDBExtension {
@Getter
private final BanyanDB.IndexRule.IndexType indexType;
+ /**
+ * A column belong to a measure's field.
+ */
+ @Getter
+ private final boolean isMeasureField;
+
/**
* @return true if this column is a part of sharding key
*/
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBModelExtension.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBModelExtension.java
index dc8b60ad9c..73b1819b29 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBModelExtension.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBModelExtension.java
@@ -37,4 +37,12 @@ public class BanyanDBModelExtension {
@Setter
private String timestampColumn;
+ /**
+ * storeIDTag indicates whether a metric stores its ID as a tag.
+ * The installer will create a virtual string ID tag with a tree index
rule.
+ */
+ @Getter
+ @Setter
+ private boolean storeIDTag;
+
}
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java
index f196f8aa5b..bd7b5d46a7 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java
@@ -93,6 +93,10 @@ public class StorageModels implements IModelManager,
ModelCreator, ModelManipula
banyanDBModelExtension.setTimestampColumn(timestampColumn);
}
+ if (aClass.isAnnotationPresent(BanyanDB.StoreIDAsTag.class)) {
+ banyanDBModelExtension.setStoreIDTag(true);
+ }
+
checker.check(storage.getModelName());
Model model = new Model(
@@ -190,7 +194,7 @@ public class StorageModels implements IModelManager,
ModelCreator, ModelManipula
ElasticSearchExtension elasticSearchExtension = new
ElasticSearchExtension(
elasticSearchAnalyzer == null ? null :
elasticSearchAnalyzer.analyzer(),
elasticSearchColumn == null ? null :
elasticSearchColumn.columnAlias(),
- keywordColumn == null ? false : true
+ keywordColumn != null
);
// BanyanDB extension
@@ -202,11 +206,14 @@ public class StorageModels implements IModelManager,
ModelCreator, ModelManipula
BanyanDB.NoIndexing.class);
final BanyanDB.IndexRule banyanDBIndexRule =
field.getAnnotation(
BanyanDB.IndexRule.class);
+ final BanyanDB.MeasureField banyanDBMeasureField =
field.getAnnotation(
+ BanyanDB.MeasureField.class);
BanyanDBExtension banyanDBExtension = new BanyanDBExtension(
banyanDBSeriesID == null ? -1 : banyanDBSeriesID.index(),
banyanDBGlobalIndex != null,
- banyanDBNoIndex == null && column.storageOnly(),
- banyanDBIndexRule == null ?
BanyanDB.IndexRule.IndexType.INVERTED : banyanDBIndexRule.indexType()
+ banyanDBNoIndex == null && !column.storageOnly(),
+ banyanDBIndexRule == null ?
BanyanDB.IndexRule.IndexType.INVERTED : banyanDBIndexRule.indexType(),
+ banyanDBMeasureField != null
);
final ModelColumn modelColumn = new ModelColumn(
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceRelationTraffic.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceRelationTraffic.java
index 93f267ce04..3e7eac78e2 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceRelationTraffic.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceRelationTraffic.java
@@ -29,6 +29,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -52,10 +53,12 @@ public class ZipkinServiceRelationTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = SERVICE_NAME)
+ @BanyanDB.SeriesID(index = 0)
private String serviceName;
@Setter
@Getter
@Column(columnName = REMOTE_SERVICE_NAME)
+ @BanyanDB.SeriesID(index = 1)
private String remoteServiceName;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceSpanTraffic.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceSpanTraffic.java
index 6fd873e310..ebb1c9bf23 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceSpanTraffic.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceSpanTraffic.java
@@ -30,6 +30,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -54,10 +55,12 @@ public class ZipkinServiceSpanTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = SERVICE_NAME)
+ @BanyanDB.SeriesID(index = 0)
private String serviceName;
@Setter
@Getter
@Column(columnName = SPAN_NAME)
+ @BanyanDB.SeriesID(index = 1)
private String spanName = Const.EMPTY_STRING;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceTraffic.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceTraffic.java
index 098790919b..9d962936bf 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceTraffic.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinServiceTraffic.java
@@ -30,6 +30,7 @@ import
org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.source.DefaultScopeDefine;
import org.apache.skywalking.oap.server.core.storage.ShardingAlgorithm;
import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SQLDatabase;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
@@ -51,6 +52,7 @@ public class ZipkinServiceTraffic extends Metrics {
@Setter
@Getter
@Column(columnName = SERVICE_NAME)
+ @BanyanDB.SeriesID(index = 0)
private String serviceName = Const.EMPTY_STRING;
@Override
diff --git
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java
index 959476867f..0d58266a31 100644
---
a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java
+++
b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java
@@ -83,10 +83,12 @@ public class ZipkinSpanRecord extends Record {
@Getter
@Column(columnName = TRACE_ID)
@SQLDatabase.AdditionalEntity(additionalTables = {ADDITIONAL_QUERY_TABLE},
reserveOriginalColumns = true)
+ @BanyanDB.SeriesID(index = 0)
private String traceId;
@Setter
@Getter
@Column(columnName = SPAN_ID)
+ @BanyanDB.SeriesID(index = 1)
private String spanId;
@Setter
@Getter
@@ -115,7 +117,6 @@ public class ZipkinSpanRecord extends Record {
@Setter
@Getter
@Column(columnName = LOCAL_ENDPOINT_SERVICE_NAME)
- @BanyanDB.SeriesID(index = 0)
private String localEndpointServiceName;
@Setter
@Getter
diff --git
a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/model/ModelColumnTest.java
b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/model/ModelColumnTest.java
index 51c6cb9e62..5cd16665fd 100644
---
a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/model/ModelColumnTest.java
+++
b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/model/ModelColumnTest.java
@@ -32,7 +32,7 @@ public class ModelColumnTest {
new SQLDatabaseExtension(),
new ElasticSearchExtension(
ElasticSearch.MatchQuery.AnalyzerType.OAP_ANALYZER, "abc", false),
- new BanyanDBExtension(-1, false,
true, BanyanDB.IndexRule.IndexType.INVERTED)
+ new BanyanDBExtension(-1, false,
true, BanyanDB.IndexRule.IndexType.INVERTED, false)
);
Assert.assertEquals(true, column.isStorageOnly());
Assert.assertEquals("abc", column.getColumnName().getName());
@@ -41,7 +41,7 @@ public class ModelColumnTest {
false, false, true, 200,
new SQLDatabaseExtension(),
new
ElasticSearchExtension(ElasticSearch.MatchQuery.AnalyzerType.OAP_ANALYZER,
"abc", false),
- new BanyanDBExtension(-1, false, true,
BanyanDB.IndexRule.IndexType.INVERTED)
+ new BanyanDBExtension(-1, false, true,
BanyanDB.IndexRule.IndexType.INVERTED, false)
);
Assert.assertEquals(true, column.isStorageOnly());
Assert.assertEquals("abc", column.getColumnName().getName());
@@ -51,7 +51,7 @@ public class ModelColumnTest {
false, false, true, 200,
new SQLDatabaseExtension(),
new
ElasticSearchExtension(ElasticSearch.MatchQuery.AnalyzerType.OAP_ANALYZER,
"abc", false),
- new BanyanDBExtension(-1, false, true,
BanyanDB.IndexRule.IndexType.INVERTED)
+ new BanyanDBExtension(-1, false, true,
BanyanDB.IndexRule.IndexType.INVERTED, false)
);
Assert.assertEquals(false, column.isStorageOnly());
Assert.assertEquals("abc", column.getColumnName().getName());
@@ -64,7 +64,7 @@ public class ModelColumnTest {
new SQLDatabaseExtension(),
new ElasticSearchExtension(
ElasticSearch.MatchQuery.AnalyzerType.OAP_ANALYZER, "abc", false),
- new BanyanDBExtension(-1, false,
true, BanyanDB.IndexRule.IndexType.INVERTED)
+ new BanyanDBExtension(-1, false,
true, BanyanDB.IndexRule.IndexType.INVERTED, false)
);
}
@@ -75,7 +75,7 @@ public class ModelColumnTest {
new SQLDatabaseExtension(),
new ElasticSearchExtension(
ElasticSearch.MatchQuery.AnalyzerType.OAP_ANALYZER, "abc", false),
- new BanyanDBExtension(-1, false,
true, BanyanDB.IndexRule.IndexType.INVERTED)
+ new BanyanDBExtension(-1, false,
true, BanyanDB.IndexRule.IndexType.INVERTED, false)
);
}
}
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java
index de32a409a2..270fca5135 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java
@@ -40,6 +40,9 @@ import
org.apache.skywalking.oap.server.storage.plugin.banyandb.util.ByteUtil;
import java.util.List;
public class BanyanDBConverter {
+
+ public static final String ID = "id";
+
public static class StorageToStream implements Convert2Entity {
private final MetadataRegistry.Schema schema;
private final RowEntity rowEntity;
@@ -156,7 +159,7 @@ public class BanyanDBConverter {
public void acceptID(String id) {
try {
- this.measureWrite.setID(id);
+ this.measureWrite.tag(ID, TagAndValue.stringTagValue(id));
} catch (BanyanDBException ex) {
log.error("fail to add ID tag", ex);
}
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java
index 2d40e2d6c9..b765dd4f63 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java
@@ -69,7 +69,7 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
return true;
}
- throw new IllegalStateException("inconsistent state");
+ throw new IllegalStateException("inconsistent state:" + metadata);
} catch (BanyanDBException ex) {
throw new StorageException("fail to check existence", ex);
}
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java
index c3ada1e898..5756bd710e 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java
@@ -45,6 +45,7 @@ import lombok.NoArgsConstructor;
import lombok.RequiredArgsConstructor;
import lombok.Setter;
import lombok.Singular;
+import lombok.ToString;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.v1.client.BanyanDBClient;
import
org.apache.skywalking.banyandb.v1.client.grpc.exception.BanyanDBException;
@@ -64,6 +65,7 @@ import
org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
import org.apache.skywalking.oap.server.core.analysis.record.Record;
import org.apache.skywalking.oap.server.core.config.ConfigService;
import org.apache.skywalking.oap.server.core.query.enumeration.Step;
+import org.apache.skywalking.oap.server.core.storage.StorageException;
import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDB;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import
org.apache.skywalking.oap.server.core.storage.annotation.ValueColumnMetadata;
@@ -88,12 +90,15 @@ public enum MetadataRegistry {
.collect(Collectors.toMap(modelColumn ->
modelColumn.getColumnName().getStorageName(), Function.identity()));
// parse and set sharding keys
List<String> shardingColumns = parseEntityNames(modelColumnMap);
+ if (shardingColumns.isEmpty()) {
+ throw new IllegalStateException("sharding keys of model[stream." +
model.getName() + "] must not be empty");
+ }
// parse tag metadata
// this can be used to build both
// 1) a list of TagFamilySpec,
// 2) a list of IndexRule,
- List<TagMetadata> tags = parseTagMetadata(model, schemaBuilder);
- List<TagFamilySpec> tagFamilySpecs =
schemaMetadata.extractTagFamilySpec(tags);
+ List<TagMetadata> tags = parseTagMetadata(model, schemaBuilder,
shardingColumns);
+ List<TagFamilySpec> tagFamilySpecs =
schemaMetadata.extractTagFamilySpec(tags, false);
// iterate over tagFamilySpecs to save tag names
for (final TagFamilySpec tagFamilySpec : tagFamilySpecs) {
for (final TagFamilySpec.TagSpec tagSpec :
tagFamilySpec.tagSpecs()) {
@@ -112,9 +117,6 @@ public enum MetadataRegistry {
.collect(Collectors.toList());
final Stream.Builder builder =
Stream.create(schemaMetadata.getGroup(), schemaMetadata.name());
- if (shardingColumns.isEmpty()) {
- throw new IllegalStateException("sharding keys of model[stream." +
model.getName() + "] must not be empty");
- }
builder.setEntityRelativeTags(shardingColumns);
builder.addTagFamilies(tagFamilySpecs);
builder.addIndexes(indexRules);
@@ -122,45 +124,49 @@ public enum MetadataRegistry {
return builder.build();
}
- public Measure registerMeasureModel(Model model, BanyanDBStorageConfig
config, ConfigService configService) {
+ public Measure registerMeasureModel(Model model, BanyanDBStorageConfig
config, ConfigService configService) throws StorageException {
final SchemaMetadata schemaMetadata = parseMetadata(model, config,
configService);
Schema.SchemaBuilder schemaBuilder =
Schema.builder().metadata(schemaMetadata);
Map<String, ModelColumn> modelColumnMap = model.getColumns().stream()
.collect(Collectors.toMap(modelColumn ->
modelColumn.getColumnName().getStorageName(), Function.identity()));
// parse and set sharding keys
List<String> shardingColumns = parseEntityNames(modelColumnMap);
+ if (shardingColumns.isEmpty()) {
+ throw new StorageException("model " + model.getName() + " doesn't
contain series id");
+ }
// parse tag metadata
// this can be used to build both
// 1) a list of TagFamilySpec,
// 2) a list of IndexRule,
- List<TagMetadata> tags = parseTagAndFieldMetadata(model,
schemaBuilder);
- List<TagFamilySpec> tagFamilySpecs =
schemaMetadata.extractTagFamilySpec(tags);
+ MeasureMetadata tagsAndFields = parseTagAndFieldMetadata(model,
schemaBuilder, shardingColumns);
+ List<TagFamilySpec> tagFamilySpecs =
schemaMetadata.extractTagFamilySpec(tagsAndFields.tags,
model.getBanyanDBModelExtension().isStoreIDTag());
// iterate over tagFamilySpecs to save tag names
for (final TagFamilySpec tagFamilySpec : tagFamilySpecs) {
for (final TagFamilySpec.TagSpec tagSpec :
tagFamilySpec.tagSpecs()) {
schemaBuilder.tag(tagSpec.getTagName());
}
}
- List<IndexRule> indexRules = tags.stream()
+ List<IndexRule> indexRules = tagsAndFields.tags.stream()
.map(TagMetadata::getIndexRule)
.filter(Objects::nonNull)
.collect(Collectors.toList());
+ if (model.getBanyanDBModelExtension().isStoreIDTag()) {
+ indexRules.add(IndexRule.create(BanyanDBConverter.ID,
IndexRule.IndexType.TREE, IndexRule.IndexLocation.SERIES));
+ }
+
final Measure.Builder builder =
Measure.create(schemaMetadata.getGroup(), schemaMetadata.name(),
downSamplingDuration(model.getDownsampling()));
- if (shardingColumns.isEmpty()) {
- // if shardingKeys is empty, for measure, we can use ID as a
single sharding key.
- builder.setEntityRelativeTags(Measure.ID);
- } else {
- builder.setEntityRelativeTags(shardingColumns);
- }
+ builder.setEntityRelativeTags(shardingColumns);
builder.addTagFamilies(tagFamilySpecs);
- builder.addIndexes(indexRules);
+ if (!indexRules.isEmpty()) {
+ builder.addIndexes(indexRules);
+ }
// parse and set field
- Optional<ValueColumnMetadata.ValueColumn> valueColumnOpt =
ValueColumnMetadata.INSTANCE
- .readValueColumnDefinition(model.getName());
- valueColumnOpt.ifPresent(valueColumn ->
builder.addField(parseFieldSpec(modelColumnMap.get(valueColumn.getValueCName()),
valueColumn)));
- valueColumnOpt.ifPresent(valueColumn ->
schemaBuilder.field(valueColumn.getValueCName()));
+ for (Measure.FieldSpec field : tagsAndFields.fields) {
+ builder.addField(field);
+ schemaBuilder.field(field.getName());
+ }
registry.put(schemaMetadata.name(), schemaBuilder.build());
return builder.build();
}
@@ -208,24 +214,25 @@ public enum MetadataRegistry {
return this.registry.get(SchemaMetadata.formatName(modelName,
downSampling));
}
- private Measure.FieldSpec parseFieldSpec(ModelColumn modelColumn,
ValueColumnMetadata.ValueColumn valueColumn) {
+ private Measure.FieldSpec parseFieldSpec(ModelColumn modelColumn) {
+ String colName = modelColumn.getColumnName().getStorageName();
if (String.class.equals(modelColumn.getType())) {
- return Measure.FieldSpec.newIntField(valueColumn.getValueCName())
+ return Measure.FieldSpec.newIntField(colName)
.compressWithZSTD()
.build();
} else if (long.class.equals(modelColumn.getType()) ||
int.class.equals(modelColumn.getType())) {
- return Measure.FieldSpec.newIntField(valueColumn.getValueCName())
+ return Measure.FieldSpec.newIntField(colName)
.compressWithZSTD()
.encodeWithGorilla()
.build();
- } else if
(StorageDataComplexObject.class.isAssignableFrom(modelColumn.getType())) {
- return
Measure.FieldSpec.newStringField(valueColumn.getValueCName())
+ } else if
(StorageDataComplexObject.class.isAssignableFrom(modelColumn.getType()) ||
JsonObject.class.equals(modelColumn.getType())) {
+ return Measure.FieldSpec.newStringField(colName)
.compressWithZSTD()
.build();
} else if (double.class.equals(modelColumn.getType())) {
// TODO: natively support double/float in BanyanDB
log.warn("Double is stored as binary");
- return
Measure.FieldSpec.newBinaryField(valueColumn.getValueCName())
+ return Measure.FieldSpec.newBinaryField(colName)
.compressWithZSTD()
.build();
} else {
@@ -293,7 +300,7 @@ public enum MetadataRegistry {
*
* @since 9.4.0 Skip {@link Record#TIME_BUCKET}
*/
- List<TagMetadata> parseTagMetadata(Model model, Schema.SchemaBuilder
builder) {
+ List<TagMetadata> parseTagMetadata(Model model, Schema.SchemaBuilder
builder, List<String> shardingColumns) {
List<TagMetadata> tagMetadataList = new ArrayList<>();
for (final ModelColumn col : model.getColumns()) {
final String columnStorageName =
col.getColumnName().getStorageName();
@@ -302,7 +309,8 @@ public enum MetadataRegistry {
}
final TagFamilySpec.TagSpec tagSpec = parseTagSpec(col);
builder.spec(columnStorageName, new ColumnSpec(ColumnType.TAG,
col.getType()));
- if (col.shouldIndex()) {
+ String colName = col.getColumnName().getStorageName();
+ if (!shardingColumns.contains(colName) &&
col.getBanyanDBExtension().shouldIndex()) {
// build indexRule
IndexRule indexRule = parseIndexRule(tagSpec.getTagName(),
col);
tagMetadataList.add(new TagMetadata(indexRule, tagSpec));
@@ -314,6 +322,14 @@ public enum MetadataRegistry {
return tagMetadataList;
}
+ @Builder
+ private static class MeasureMetadata {
+ @Singular
+ private final List<TagMetadata> tags;
+ @Singular
+ private final List<Measure.FieldSpec> fields;
+ }
+
/**
* Parse tags and fields' metadata for {@link Measure}.
* For field whose dataType is not {@link Column.ValueDataType#NOT_VALUE},
@@ -321,32 +337,28 @@ public enum MetadataRegistry {
*
* @since 9.4.0 Skip {@link Metrics#TIME_BUCKET}
*/
- List<TagMetadata> parseTagAndFieldMetadata(Model model,
Schema.SchemaBuilder builder) {
- List<TagMetadata> tagMetadataList = new ArrayList<>();
+ MeasureMetadata parseTagAndFieldMetadata(Model model, Schema.SchemaBuilder
builder, List<String> shardingColumns) {
// skip metric
Optional<ValueColumnMetadata.ValueColumn> valueColumnOpt =
ValueColumnMetadata.INSTANCE
.readValueColumnDefinition(model.getName());
+ MeasureMetadata.MeasureMetadataBuilder result =
MeasureMetadata.builder();
for (final ModelColumn col : model.getColumns()) {
final String columnStorageName =
col.getColumnName().getStorageName();
if (columnStorageName.equals(Metrics.TIME_BUCKET)) {
continue;
}
- if (valueColumnOpt.isPresent() &&
valueColumnOpt.get().getValueCName().equals(columnStorageName)) {
+ if (col.getBanyanDBExtension().isMeasureField()) {
builder.spec(columnStorageName, new
ColumnSpec(ColumnType.FIELD, col.getType()));
+ result.field(parseFieldSpec(col));
continue;
}
final TagFamilySpec.TagSpec tagSpec = parseTagSpec(col);
builder.spec(columnStorageName, new ColumnSpec(ColumnType.TAG,
col.getType()));
- if (col.shouldIndex()) {
- // build indexRule
- IndexRule indexRule = parseIndexRule(tagSpec.getTagName(),
col);
- tagMetadataList.add(new TagMetadata(indexRule, tagSpec));
- } else {
- tagMetadataList.add(new TagMetadata(null, tagSpec));
- }
+ String colName = col.getColumnName().getStorageName();
+ result.tag(new TagMetadata(!shardingColumns.contains(colName) &&
col.getBanyanDBExtension().shouldIndex() ? parseIndexRule(tagSpec.getTagName(),
col) : null, tagSpec));
}
- return tagMetadataList;
+ return result.build();
}
/**
@@ -462,6 +474,7 @@ public enum MetadataRegistry {
@RequiredArgsConstructor
@Data
+ @ToString
public static class SchemaMetadata {
private final String group;
/**
@@ -511,18 +524,18 @@ public enum MetadataRegistry {
}
}
- private List<TagFamilySpec> extractTagFamilySpec(List<TagMetadata>
tagMetadataList) {
+ private List<TagFamilySpec> extractTagFamilySpec(List<TagMetadata>
tagMetadataList, boolean shouldAddID) {
+ final String indexFamily = SchemaMetadata.this.indexFamily();
+ final String nonIndexFamily = SchemaMetadata.this.nonIndexFamily();
Map<String, List<TagMetadata>> tagMetadataMap =
tagMetadataList.stream()
- .collect(Collectors.groupingBy(tagMetadata ->
tagMetadata.isIndex() ? SchemaMetadata.this.indexFamily() :
SchemaMetadata.this.nonIndexFamily()));
+ .collect(Collectors.groupingBy(tagMetadata ->
tagMetadata.isIndex() ? indexFamily : nonIndexFamily));
final List<TagFamilySpec> tagFamilySpecs = new
ArrayList<>(tagMetadataMap.size());
for (final Map.Entry<String, List<TagMetadata>> entry :
tagMetadataMap.entrySet()) {
final TagFamilySpec.Builder b =
TagFamilySpec.create(entry.getKey())
.addTagSpecs(entry.getValue().stream().map(TagMetadata::getTagSpec).collect(Collectors.toList()));
- if (this.getKind() == Kind.MEASURE &&
entry.getKey().equals(this.indexFamily())) {
- // append measure ID, but it should not generate an index
in the client side.
- // BanyanDB will take care of the ID index registration.
- b.addIDTagSpec();
+ if (shouldAddID && indexFamily.equals(entry.getKey())) {
+
b.addTagSpec(TagFamilySpec.TagSpec.newStringTag(BanyanDBConverter.ID));
}
tagFamilySpecs.add(b.build());
}
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBEBPFProfilingScheduleQueryDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBEBPFProfilingScheduleQueryDAO.java
index 8e3b41741a..7105d60009 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBEBPFProfilingScheduleQueryDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBEBPFProfilingScheduleQueryDAO.java
@@ -18,7 +18,8 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb.measure;
-import com.google.common.collect.ImmutableSet;
+ import com.google.common.collect.ImmutableSet;
+ import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
import org.apache.skywalking.banyandb.v1.client.DataPoint;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
@@ -35,38 +36,40 @@ import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;
-public class BanyanDBEBPFProfilingScheduleQueryDAO extends AbstractBanyanDBDAO
implements IEBPFProfilingScheduleDAO {
- private static final Set<String> TAGS =
ImmutableSet.of(EBPFProfilingScheduleRecord.START_TIME,
- EBPFProfilingScheduleRecord.TASK_ID,
- EBPFProfilingScheduleRecord.PROCESS_ID,
- EBPFProfilingScheduleRecord.END_TIME);
+@Slf4j
+ public class BanyanDBEBPFProfilingScheduleQueryDAO extends
AbstractBanyanDBDAO implements IEBPFProfilingScheduleDAO {
+ private static final Set<String> TAGS =
ImmutableSet.of(EBPFProfilingScheduleRecord.START_TIME,
+ EBPFProfilingScheduleRecord.EBPF_PROFILING_SCHEDULE_ID,
+ EBPFProfilingScheduleRecord.TASK_ID,
+ EBPFProfilingScheduleRecord.PROCESS_ID,
+ EBPFProfilingScheduleRecord.END_TIME);
- public BanyanDBEBPFProfilingScheduleQueryDAO(BanyanDBStorageClient client)
{
- super(client);
- }
+ public BanyanDBEBPFProfilingScheduleQueryDAO(BanyanDBStorageClient
client) {
+ super(client);
+ }
- @Override
- public List<EBPFProfilingSchedule> querySchedules(String taskId) throws
IOException {
- MeasureQueryResponse resp =
query(EBPFProfilingScheduleRecord.INDEX_NAME,
- TAGS,
- Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
- @Override
- protected void apply(MeasureQuery query) {
- query.and(eq(EBPFProfilingScheduleRecord.TASK_ID,
taskId));
- query.setOrderBy(new
AbstractQuery.OrderBy(EBPFProfilingScheduleRecord.START_TIME,
AbstractQuery.Sort.DESC));
- }
- });
+ @Override
+ public List<EBPFProfilingSchedule> querySchedules(String taskId) throws
IOException {
+ MeasureQueryResponse resp =
query(EBPFProfilingScheduleRecord.INDEX_NAME,
+ TAGS,
+ Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
+ @Override
+ protected void apply(MeasureQuery query) {
+ query.and(eq(EBPFProfilingScheduleRecord.TASK_ID,
taskId));
+ query.setOrderBy(new
AbstractQuery.OrderBy(EBPFProfilingScheduleRecord.START_TIME,
AbstractQuery.Sort.DESC));
+ }
+ });
- return
resp.getDataPoints().stream().map(this::buildEBPFProfilingSchedule).collect(Collectors.toList());
- }
+ return
resp.getDataPoints().stream().map(this::buildEBPFProfilingSchedule).collect(Collectors.toList());
+ }
- private EBPFProfilingSchedule buildEBPFProfilingSchedule(DataPoint
dataPoint) {
- final EBPFProfilingSchedule schedule = new EBPFProfilingSchedule();
- schedule.setScheduleId(dataPoint.getId());
-
schedule.setTaskId(dataPoint.getTagValue(EBPFProfilingScheduleRecord.TASK_ID));
-
schedule.setProcessId(dataPoint.getTagValue(EBPFProfilingScheduleRecord.PROCESS_ID));
- schedule.setStartTime(((Number)
dataPoint.getTagValue(EBPFProfilingScheduleRecord.START_TIME)).longValue());
- schedule.setEndTime(((Number)
dataPoint.getTagValue(EBPFProfilingScheduleRecord.END_TIME)).longValue());
- return schedule;
- }
-}
+ private EBPFProfilingSchedule buildEBPFProfilingSchedule(DataPoint
dataPoint) {
+ final EBPFProfilingSchedule schedule = new EBPFProfilingSchedule();
+
schedule.setScheduleId(dataPoint.getTagValue(EBPFProfilingScheduleRecord.EBPF_PROFILING_SCHEDULE_ID));
+
schedule.setTaskId(dataPoint.getTagValue(EBPFProfilingScheduleRecord.TASK_ID));
+
schedule.setProcessId(dataPoint.getTagValue(EBPFProfilingScheduleRecord.PROCESS_ID));
+ schedule.setStartTime(((Number)
dataPoint.getTagValue(EBPFProfilingScheduleRecord.START_TIME)).longValue());
+ schedule.setEndTime(((Number)
dataPoint.getTagValue(EBPFProfilingScheduleRecord.END_TIME)).longValue());
+ return schedule;
+ }
+ }
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetadataQueryDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetadataQueryDAO.java
index d10b92a22f..3231a506fc 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetadataQueryDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetadataQueryDAO.java
@@ -18,14 +18,15 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb.measure;
-import com.google.common.base.Strings;
import com.google.common.collect.ImmutableSet;
import com.google.gson.Gson;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
+import org.apache.commons.lang3.StringUtils;
import org.apache.skywalking.banyandb.v1.client.DataPoint;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
import org.apache.skywalking.banyandb.v1.client.MeasureQueryResponse;
+import org.apache.skywalking.oap.server.core.analysis.DownSampling;
import org.apache.skywalking.oap.server.core.analysis.IDManager;
import org.apache.skywalking.oap.server.core.analysis.Layer;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
@@ -44,7 +45,9 @@ import
org.apache.skywalking.oap.server.core.query.type.Service;
import org.apache.skywalking.oap.server.core.query.type.ServiceInstance;
import org.apache.skywalking.oap.server.core.storage.query.IMetadataQueryDAO;
import org.apache.skywalking.oap.server.library.util.StringUtil;
+import
org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBConverter;
import
org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageClient;
+import
org.apache.skywalking.oap.server.storage.plugin.banyandb.MetadataRegistry;
import
org.apache.skywalking.oap.server.storage.plugin.banyandb.stream.AbstractBanyanDBDAO;
import java.io.IOException;
@@ -103,9 +106,10 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
});
final List<Service> services = new ArrayList<>();
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(ServiceTraffic.INDEX_NAME,
DownSampling.Minute);
for (final DataPoint dataPoint : resp.getDataPoints()) {
- services.add(buildService(dataPoint));
+ services.add(buildService(dataPoint, schema));
}
return services;
@@ -125,9 +129,10 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
});
final List<Service> services = new ArrayList<>();
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(ServiceTraffic.INDEX_NAME,
DownSampling.Minute);
for (final DataPoint dataPoint : resp.getDataPoints()) {
- services.add(buildService(dataPoint));
+ services.add(buildService(dataPoint, schema));
}
return services;
@@ -150,9 +155,9 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
});
final List<ServiceInstance> instances = new ArrayList<>();
-
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(InstanceTraffic.INDEX_NAME,
DownSampling.Minute);
for (final DataPoint dataPoint : resp.getDataPoints()) {
- instances.add(buildInstance(dataPoint));
+ instances.add(buildInstance(dataPoint, schema));
}
return instances;
@@ -160,19 +165,19 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
@Override
public ServiceInstance getInstance(String instanceId) throws IOException {
+ IDManager.ServiceInstanceID.InstanceIDDefinition id =
IDManager.ServiceInstanceID.analysisId(instanceId);
MeasureQueryResponse resp = query(InstanceTraffic.INDEX_NAME,
INSTANCE_TRAFFIC_COMPACT_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
- if (StringUtil.isNotEmpty(instanceId)) {
- query.and(id(instanceId));
- }
+ query.and(eq(InstanceTraffic.SERVICE_ID,
id.getServiceId()))
+ .and(eq(InstanceTraffic.NAME,
id.getName()));
}
});
-
- return resp.size() > 0 ? buildInstance(resp.getDataPoints().get(0)) :
null;
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(InstanceTraffic.INDEX_NAME,
DownSampling.Minute);
+ return resp.size() > 0 ? buildInstance(resp.getDataPoints().get(0),
schema) : null;
}
@Override
@@ -190,9 +195,9 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
});
final List<Endpoint> endpoints = new ArrayList<>();
-
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(EndpointTraffic.INDEX_NAME,
DownSampling.Minute);
for (final DataPoint dataPoint : resp.getDataPoints()) {
- endpoints.add(buildEndpoint(dataPoint));
+ endpoints.add(buildEndpoint(dataPoint, schema));
}
if (StringUtil.isNotEmpty(serviceId)) {
@@ -217,9 +222,9 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
});
final List<Process> processes = new ArrayList<>();
-
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME,
DownSampling.Minute);
for (final DataPoint dataPoint : resp.getDataPoints()) {
- processes.add(buildProcess(dataPoint));
+ processes.add(buildProcess(dataPoint, schema));
}
return processes;
@@ -243,9 +248,9 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
});
final List<Process> processes = new ArrayList<>();
-
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME,
DownSampling.Minute);
for (final DataPoint dataPoint : resp.getDataPoints()) {
- processes.add(buildProcess(dataPoint));
+ processes.add(buildProcess(dataPoint, schema));
}
return processes;
@@ -265,9 +270,9 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
});
final List<Process> processes = new ArrayList<>();
-
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME,
DownSampling.Minute);
for (final DataPoint dataPoint : resp.getDataPoints()) {
- processes.add(buildProcess(dataPoint));
+ processes.add(buildProcess(dataPoint, schema));
}
return processes;
@@ -322,35 +327,38 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
@Override
protected void apply(MeasureQuery query) {
if (StringUtil.isNotEmpty(processId)) {
- query.and(id(processId));
+ query.and(eq(BanyanDBConverter.ID, processId));
}
}
});
+ MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME,
DownSampling.Minute);
- return resp.size() > 0 ? buildProcess(resp.getDataPoints().get(0)) :
null;
+ return resp.size() > 0 ? buildProcess(resp.getDataPoints().get(0),
schema) : null;
}
- private Service buildService(DataPoint dataPoint) {
+ private Service buildService(DataPoint dataPoint, MetadataRegistry.Schema
schema) {
+ final ServiceTraffic.Builder builder = new ServiceTraffic.Builder();
+ final ServiceTraffic serviceTraffic = builder.storage2Entity(new
BanyanDBConverter.StorageToMeasure(schema, dataPoint));
+ String serviceName = serviceTraffic.getName();
Service service = new Service();
- service.setId(dataPoint.getTagValue(ServiceTraffic.SERVICE_ID));
- service.setName(dataPoint.getTagValue(ServiceTraffic.NAME));
- service.setShortName(dataPoint.getTagValue(ServiceTraffic.SHORT_NAME));
- service.setGroup(dataPoint.getTagValue(ServiceTraffic.GROUP));
- service.getLayers().add(Layer.valueOf(((Number)
dataPoint.getTagValue(ServiceTraffic.LAYER)).intValue()).name());
+ service.setId(serviceTraffic.getServiceId());
+ service.setName(serviceName);
+ service.setShortName(serviceTraffic.getShortName());
+ service.setGroup(serviceTraffic.getGroup());
+ service.getLayers().add(serviceTraffic.getLayer().name());
return service;
}
- private ServiceInstance buildInstance(DataPoint dataPoint) {
+ private ServiceInstance buildInstance(DataPoint dataPoint,
MetadataRegistry.Schema schema) {
+ final InstanceTraffic instanceTraffic =
+ new InstanceTraffic.Builder().storage2Entity(new
BanyanDBConverter.StorageToMeasure(schema, dataPoint));
+
ServiceInstance serviceInstance = new ServiceInstance();
- serviceInstance.setId(dataPoint.getId());
- serviceInstance.setName(dataPoint.getTagValue(InstanceTraffic.NAME));
- serviceInstance.setInstanceUUID(dataPoint.getId());
-
- final String propString =
dataPoint.getTagValue(InstanceTraffic.PROPERTIES);
- JsonObject properties = null;
- if (StringUtil.isNotEmpty(propString)) {
- properties = GSON.fromJson(propString, JsonObject.class);
- }
+ serviceInstance.setId(instanceTraffic.id().build());
+ serviceInstance.setName(instanceTraffic.getName());
+ serviceInstance.setInstanceUUID(serviceInstance.getId());
+
+ JsonObject properties = instanceTraffic.getProperties();
if (properties != null) {
for (Map.Entry<String, JsonElement> property :
properties.entrySet()) {
String key = property.getKey();
@@ -364,44 +372,46 @@ public class BanyanDBMetadataQueryDAO extends
AbstractBanyanDBDAO implements IMe
} else {
serviceInstance.setLanguage(Language.UNKNOWN);
}
-
return serviceInstance;
}
- private Endpoint buildEndpoint(DataPoint dataPoint) {
+ private Endpoint buildEndpoint(DataPoint dataPoint,
MetadataRegistry.Schema schema) {
+ final EndpointTraffic endpointTraffic =
+ new EndpointTraffic.Builder().storage2Entity(new
BanyanDBConverter.StorageToMeasure(schema, dataPoint));
Endpoint endpoint = new Endpoint();
- endpoint.setId(dataPoint.getId());
- endpoint.setName(dataPoint.getTagValue(EndpointTraffic.NAME));
+ endpoint.setId(endpointTraffic.id().build());
+ endpoint.setName(endpointTraffic.getName());
return endpoint;
}
- private Process buildProcess(DataPoint dataPoint) {
- Process process = new Process();
+ private Process buildProcess(DataPoint dataPoint, MetadataRegistry.Schema
schema) {
+ final ProcessTraffic processTraffic =
+ new ProcessTraffic.Builder().storage2Entity(new
BanyanDBConverter.StorageToMeasure(schema, dataPoint));
- process.setId(dataPoint.getId());
- process.setName(dataPoint.getTagValue(ProcessTraffic.NAME));
- String serviceId = dataPoint.getTagValue(ProcessTraffic.SERVICE_ID);
+ Process process = new Process();
+ process.setId(processTraffic.id().build());
+ process.setName(processTraffic.getName());
+ final String serviceId = processTraffic.getServiceId();
process.setServiceId(serviceId);
process.setServiceName(IDManager.ServiceID.analysisId(serviceId).getName());
- String instanceId = dataPoint.getTagValue(ProcessTraffic.INSTANCE_ID);
+ final String instanceId = processTraffic.getInstanceId();
process.setInstanceId(instanceId);
process.setInstanceName(IDManager.ServiceInstanceID.analysisId(instanceId).getName());
- process.setAgentId(dataPoint.getTagValue(ProcessTraffic.AGENT_ID));
- process.setDetectType(ProcessDetectType.valueOf(((Number)
dataPoint.getTagValue(ProcessTraffic.DETECT_TYPE)).intValue()).name());
-
process.setProfilingSupportStatus(ProfilingSupportStatus.valueOf(((Number)
dataPoint.getTagValue(ProcessTraffic.PROFILING_SUPPORT_STATUS)).intValue()).name());
+ process.setAgentId(processTraffic.getAgentId());
+
process.setDetectType(ProcessDetectType.valueOf(processTraffic.getDetectType()).name());
+
process.setProfilingSupportStatus(ProfilingSupportStatus.valueOf(processTraffic.getProfilingSupportStatus()).name());
- String propString = dataPoint.getTagValue(ProcessTraffic.PROPERTIES);
- if (!Strings.isNullOrEmpty(propString)) {
- JsonObject properties = GSON.fromJson(propString,
JsonObject.class);
+ JsonObject properties = processTraffic.getProperties();
+ if (properties != null) {
for (Map.Entry<String, JsonElement> property :
properties.entrySet()) {
String key = property.getKey();
String value = property.getValue().getAsString();
process.getAttributes().add(new Attribute(key, value));
}
}
- String labelJson = dataPoint.getTagValue(ProcessTraffic.LABELS_JSON);
- if (!Strings.isNullOrEmpty(labelJson)) {
- List<String> labels = GSON.<List<String>>fromJson(labelJson,
ArrayList.class);
+ final String labelsJson = processTraffic.getLabelsJson();
+ if (StringUtils.isNotEmpty(labelsJson)) {
+ final List<String> labels =
GSON.<List<String>>fromJson(labelsJson, ArrayList.class);
process.getLabels().addAll(labels);
}
return process;
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsDAO.java
index 758098eb63..767da1cb92 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsDAO.java
@@ -18,15 +18,20 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb.measure;
+import com.google.common.base.Preconditions;
import lombok.extern.slf4j.Slf4j;
+import org.apache.logging.log4j.util.Strings;
import org.apache.skywalking.banyandb.v1.client.DataPoint;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
import org.apache.skywalking.banyandb.v1.client.MeasureQueryResponse;
import org.apache.skywalking.banyandb.v1.client.MeasureWrite;
+import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
import org.apache.skywalking.oap.server.core.storage.IMetricsDAO;
import org.apache.skywalking.oap.server.core.storage.SessionCacheCallback;
+import org.apache.skywalking.oap.server.core.storage.StorageID;
+import org.apache.skywalking.oap.server.core.storage.model.BanyanDBExtension;
import org.apache.skywalking.oap.server.core.storage.model.Model;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.library.client.request.InsertRequest;
@@ -39,7 +44,13 @@ import
org.apache.skywalking.oap.server.storage.plugin.banyandb.stream.AbstractB
import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+
+import static
org.apache.skywalking.oap.server.core.analysis.metrics.Metrics.ENTITY_ID;
+import static
org.apache.skywalking.oap.server.core.storage.StorageData.TIME_BUCKET;
@Slf4j
public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements
IMetricsDAO {
@@ -57,16 +68,54 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO
implements IMetricsD
if (schema == null) {
throw new IOException(model.getName() + " is not registered");
}
+ String tc = model.getBanyanDBModelExtension().getTimestampColumn();
+ final String tsCol = Strings.isBlank(tc) ? TIME_BUCKET : tc;
+ long begin = 0L, end = 0L;
+ final Map<String, List<String>> seriesIDColumns = new HashMap<>();
+ model.getColumns().forEach(c -> {
+ BanyanDBExtension ext = c.getBanyanDBExtension();
+ if (ext == null) {
+ return;
+ }
+ if (ext.isShardingKey()) {
+ seriesIDColumns.put(c.getColumnName().getName(), new
ArrayList<>());
+ }
+ });
+ if (seriesIDColumns.isEmpty()) {
+ seriesIDColumns.put(ENTITY_ID, new ArrayList<>());
+ }
+ StringBuilder idStr = new StringBuilder();
+ for (Metrics m : metrics) {
+ AnalyticalResult result = analyze(m, tsCol, seriesIDColumns);
+
idStr.append(result.cols()).append("=").append(m.id().build()).append(",");
+ if (!result.success) {
+ continue;
+ }
+ if (begin == 0 || result.begin < begin) {
+ begin = result.begin;
+ }
+ if (end == 0 || result.end > end) {
+ end = result.end;
+ }
+ }
+ TimestampRange timestampRange = null;
+ if (begin != 0L || end != 0L) {
+ timestampRange = new TimestampRange(begin, end);
+ } else {
+ log.info("{}[{}] will scan all blocks", model.getName(), idStr);
+ }
+
List<Metrics> metricsInStorage = new ArrayList<>(metrics.size());
- MeasureQueryResponse resp = query(model.getName(), schema.getTags(),
schema.getFields(), new QueryBuilder<MeasureQuery>() {
- @Override
+ MeasureQueryResponse resp = query(model.getName(), schema.getTags(),
schema.getFields(), timestampRange, new QueryBuilder<MeasureQuery>() {
+ @Override
protected void apply(MeasureQuery query) {
- for (final Metrics missCachedMetric : metrics) {
- query.or(id(missCachedMetric.id().build()));
- }
+ seriesIDColumns.entrySet().forEach(entry -> {
+ if (!entry.getValue().isEmpty()) {
+ query.or(in(entry.getKey(), entry.getValue()));
+ }
+ });
}
});
-
if (resp.size() == 0) {
return Collections.emptyList();
}
@@ -89,7 +138,9 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO
implements IMetricsD
TimeBucket.getTimestamp(metrics.getTimeBucket(),
model.getDownsampling())); // timestamp
final BanyanDBConverter.MeasureToStorage toStorage = new
BanyanDBConverter.MeasureToStorage(schema, measureWrite);
storageBuilder.entity2Storage(metrics, toStorage);
- toStorage.acceptID(metrics.id().build());
+ if (model.getBanyanDBModelExtension().isStoreIDTag()) {
+ toStorage.acceptID(metrics.id().build());
+ }
return new BanyanDBMeasureInsertRequest(toStorage.obtain(), callback);
}
@@ -105,7 +156,63 @@ public class BanyanDBMetricsDAO extends
AbstractBanyanDBDAO implements IMetricsD
TimeBucket.getTimestamp(metrics.getTimeBucket(),
model.getDownsampling())); // timestamp
final BanyanDBConverter.MeasureToStorage toStorage = new
BanyanDBConverter.MeasureToStorage(schema, measureWrite);
storageBuilder.entity2Storage(metrics, toStorage);
- toStorage.acceptID(metrics.id().build());
+ if (model.getBanyanDBModelExtension().isStoreIDTag()) {
+ toStorage.acceptID(metrics.id().build());
+ }
return new BanyanDBMeasureUpdateRequest(toStorage.obtain());
}
+
+ private static class AnalyticalResult {
+ private boolean success;
+ private List<String[]> cols = new ArrayList<>();
+ private long begin;
+ private long end;
+
+ private String cols() {
+ StringBuilder b = new StringBuilder();
+ for (String[] col : this.cols) {
+ for (String c : col) {
+ b.append(c).append(",");
+ }
+ b.append(" ");
+ }
+ return b.toString();
+ }
+ }
+
+ private AnalyticalResult analyze(Metrics m, String tsCol, Map<String,
List<String>> seriesIDColumns) {
+ StorageID id = m.id();
+ List<StorageID.Fragment> fragments = id.read();
+ AnalyticalResult result = new AnalyticalResult();
+ for (StorageID.Fragment f : fragments) {
+ Optional<String[]> cols = f.getName();
+ if (cols.isPresent()) {
+ result.cols.add(cols.get());
+ for (String col : cols.get()) {
+ if (tsCol.equals(col)) {
+ long timeBucket = (long) f.getValue();
+ long epoch = TimeBucket.getTimestamp(timeBucket);
+ if (result.begin == 0 || epoch < result.begin) {
+ result.begin = epoch;
+ }
+ if (result.end == 0 || epoch > result.end) {
+ result.end = epoch;
+ }
+ } else if (seriesIDColumns.containsKey(col)) {
+
Preconditions.checkState(f.getType().equals(String.class));
+ seriesIDColumns.get(col).add((String) f.getValue());
+ } else {
+ log.error("col [{}] in fragment [{}] in id [{}] is not
ts or seriesID", col, f, id.build());
+ return result;
+ }
+ }
+ } else {
+ log.error("fragment [{}] in id [{}] doesn't contains cols", f,
id.build());
+ return result;
+ }
+ }
+ result.success = true;
+ return result;
+ }
+
}
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsQueryDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsQueryDAO.java
index ee306052c8..047f89a9c9 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsQueryDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/measure/BanyanDBMetricsQueryDAO.java
@@ -32,6 +32,7 @@ import org.apache.skywalking.banyandb.v1.client.DataPoint;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
import org.apache.skywalking.banyandb.v1.client.MeasureQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
+import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.analysis.metrics.DataTable;
import org.apache.skywalking.oap.server.core.analysis.metrics.HistogramMetrics;
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
@@ -113,33 +114,27 @@ public class BanyanDBMetricsQueryDAO extends
AbstractBanyanDBDAO implements IMet
}
final String entityID = condition.getEntity().buildId();
- Map<String, DataPoint> idMap = queryByEntityID(schema,
valueColumnName, duration, entityID);
+ Map<Long, DataPoint> idMap = queryByEntityID(schema, valueColumnName,
duration, entityID);
- List<String> ids = extractMeasureIDs(duration, entityID);
+ List<PointOfTime> tsPoints = duration.assembleDurationPoints();
MetricsValues metricsValues = new MetricsValues();
- if (!idMap.isEmpty()) {
- // Label is null, because in readMetricsValues, no label parameter.
- IntValues intValues = metricsValues.getValues();
- for (String id : ids) {
- KVInt kvInt = new KVInt();
- kvInt.setId(id);
- kvInt.setValue(0);
- if (idMap.containsKey(id)) {
- DataPoint dataPoint = idMap.get(id);
- kvInt.setValue(extractFieldValue(schema, valueColumnName,
dataPoint));
- } else {
-
kvInt.setValue(ValueColumnMetadata.INSTANCE.getDefaultValue(condition.getName()));
- }
- intValues.addKVInt(kvInt);
+ // Label is null, because in readMetricsValues, no label parameter.
+ IntValues intValues = metricsValues.getValues();
+ for (PointOfTime ts : tsPoints) {
+ String id = ts.id(entityID);
+ KVInt kvInt = new KVInt();
+ kvInt.setId(id);
+ kvInt.setValue(0);
+ if (idMap.containsKey(ts.getPoint())) {
+ DataPoint dataPoint = idMap.get(ts.getPoint());
+ kvInt.setValue(extractFieldValue(schema, valueColumnName,
dataPoint));
+ } else {
+
kvInt.setValue(ValueColumnMetadata.INSTANCE.getDefaultValue(condition.getName()));
}
+ intValues.addKVInt(kvInt);
}
- metricsValues.setValues(
- Util.sortValues(
- metricsValues.getValues(), ids,
ValueColumnMetadata.INSTANCE.getDefaultValue(condition.getName()))
- );
-
return metricsValues;
}
@@ -157,16 +152,22 @@ public class BanyanDBMetricsQueryDAO extends
AbstractBanyanDBDAO implements IMet
@Override
public List<MetricsValues> readLabeledMetricsValues(MetricsCondition
condition, String valueColumnName, List<String> labels, Duration duration)
throws IOException {
- Map<String, DataPoint> idMap = queryByEntityID(condition,
valueColumnName, duration);
+ Map<Long, DataPoint> idMap = queryByEntityID(condition,
valueColumnName, duration);
- List<String> ids = extractMeasureIDs(duration,
condition.getEntity().buildId());
+ List<PointOfTime> tsPoints = duration.assembleDurationPoints();
+ String entityID = condition.getEntity().buildId();
+ List<String> ids = new ArrayList<>(tsPoints.size());
Map<String, DataTable> dataTableMap = new HashMap<>(idMap.size());
- for (final Map.Entry<String, DataPoint> entry : idMap.entrySet()) {
- dataTableMap.put(
- entry.getKey(),
- new
DataTable(entry.getValue().getFieldValue(valueColumnName))
- );
+ for (PointOfTime ts : tsPoints) {
+ String id = ts.id(entityID);
+ ids.add(id);
+ if (idMap.containsKey(ts.getPoint())) {
+ dataTableMap.put(
+ id,
+ new
DataTable(idMap.get(ts.getPoint()).getFieldValue(valueColumnName))
+ );
+ }
}
return Util.sortValues(
@@ -178,18 +179,22 @@ public class BanyanDBMetricsQueryDAO extends
AbstractBanyanDBDAO implements IMet
@Override
public HeatMap readHeatMap(MetricsCondition condition, String
valueColumnName, Duration duration) throws IOException {
- Map<String, DataPoint> idMap = queryByEntityID(condition,
valueColumnName, duration);
+ Map<Long, DataPoint> idMap = queryByEntityID(condition,
valueColumnName, duration);
HeatMap heatMap = new HeatMap();
if (idMap.isEmpty()) {
return heatMap;
}
- List<String> ids = extractMeasureIDs(duration,
condition.getEntity().buildId());
+ List<PointOfTime> tsPoints = duration.assembleDurationPoints();
+ String entityID = condition.getEntity().buildId();
+ List<String> ids = new ArrayList<>(tsPoints.size());
final int defaultValue =
ValueColumnMetadata.INSTANCE.getDefaultValue(condition.getName());
- for (String id : ids) {
- DataPoint dataPoint = idMap.get(id);
+ for (PointOfTime ts : tsPoints) {
+ String id = ts.id(entityID);
+ ids.add(id);
+ DataPoint dataPoint = idMap.get(ts.getPoint());
if (dataPoint != null) {
String value =
dataPoint.getFieldValue(HistogramMetrics.DATASET);
heatMap.buildColumn(id, value, defaultValue);
@@ -201,17 +206,7 @@ public class BanyanDBMetricsQueryDAO extends
AbstractBanyanDBDAO implements IMet
return heatMap;
}
- private List<String> extractMeasureIDs(Duration duration, String entityID)
{
- final List<PointOfTime> pointOfTimes =
duration.assembleDurationPoints();
- List<String> ids = new ArrayList<>(pointOfTimes.size());
- pointOfTimes.forEach(pointOfTime -> {
- String id = pointOfTime.id(entityID);
- ids.add(id);
- });
- return ids;
- }
-
- private Map<String, DataPoint> queryByEntityID(final MetricsCondition
condition, String valueColumnName, Duration duration) throws IOException {
+ private Map<Long, DataPoint> queryByEntityID(final MetricsCondition
condition, String valueColumnName, Duration duration) throws IOException {
final MetadataRegistry.Schema schema =
MetadataRegistry.INSTANCE.findMetadata(condition.getName(), duration.getStep());
if (schema == null) {
throw new IOException("schema is not registered");
@@ -219,19 +214,20 @@ public class BanyanDBMetricsQueryDAO extends
AbstractBanyanDBDAO implements IMet
return queryByEntityID(schema, valueColumnName, duration,
condition.getEntity().buildId());
}
- private Map<String, DataPoint> queryByEntityID(MetadataRegistry.Schema
schema, String valueColumnName, Duration duration, String entityID) throws
IOException {
+ private Map<Long, DataPoint> queryByEntityID(MetadataRegistry.Schema
schema, String valueColumnName, Duration duration, String entityID) throws
IOException {
TimestampRange timestampRange = new
TimestampRange(duration.getStartTimestamp(), duration.getEndTimestamp());
- Map<String, DataPoint> map = new HashMap<>();
- MeasureQueryResponse resp = query(schema, Collections.emptySet(),
ImmutableSet.of(valueColumnName), timestampRange, new
QueryBuilder<MeasureQuery>() {
+ Map<Long, DataPoint> map = new HashMap<>();
+ MeasureQueryResponse resp = query(schema,
ImmutableSet.of(Metrics.ENTITY_ID), ImmutableSet.of(valueColumnName),
timestampRange, new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
query.and(eq(Metrics.ENTITY_ID, entityID));
}
});
for (final DataPoint dp : resp.getDataPoints()) {
- if (map.putIfAbsent(dp.getId(), dp) != null) {
- log.warn("duplicated data point");
+ long timeBucket = TimeBucket.getTimeBucket(dp.getTimestamp(),
schema.getMetadata().getDownSampling());
+ if (map.putIfAbsent(timeBucket, dp) != null) {
+ log.warn("duplicated data point at " + timeBucket);
}
}
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/AbstractBanyanDBDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/AbstractBanyanDBDAO.java
index 62f3facf0c..fd53ab114d 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/AbstractBanyanDBDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/AbstractBanyanDBDAO.java
@@ -29,7 +29,6 @@ import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.DownSampling;
-import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
import org.apache.skywalking.oap.server.core.storage.AbstractDAO;
import
org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageClient;
import
org.apache.skywalking.oap.server.storage.plugin.banyandb.MetadataRegistry;
@@ -146,10 +145,6 @@ public abstract class AbstractBanyanDBDAO extends
AbstractDAO<BanyanDBStorageCli
return PairQueryCondition.LongQueryCondition.ne(name, value);
}
- protected PairQueryCondition<String> id(String value) {
- return PairQueryCondition.IDQueryCondition.eq(Metrics.ID, value);
- }
-
protected AbstractQuery.OrderBy desc(String name) {
return new AbstractQuery.OrderBy(name, AbstractQuery.Sort.DESC);
}
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileTaskLogQueryDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileTaskLogQueryDAO.java
index 74038d87bc..709ce72009 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileTaskLogQueryDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileTaskLogQueryDAO.java
@@ -19,7 +19,7 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import com.google.common.collect.ImmutableSet;
-import org.apache.skywalking.banyandb.v1.client.RowEntity;
+import org.apache.skywalking.banyandb.v1.client.Element;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import
org.apache.skywalking.oap.server.core.profiling.trace.ProfileTaskLogRecord;
@@ -59,14 +59,14 @@ public class BanyanDBProfileTaskLogQueryDAO extends
AbstractBanyanDBDAO implemen
});
final LinkedList<ProfileTaskLog> tasks = new LinkedList<>();
- for (final RowEntity rowEntity : resp.getElements()) {
- tasks.add(buildProfileTaskLog(rowEntity));
+ for (final Element element : resp.getElements()) {
+ tasks.add(buildProfileTaskLog(element));
}
return tasks;
}
- private ProfileTaskLog buildProfileTaskLog(RowEntity data) {
+ private ProfileTaskLog buildProfileTaskLog(Element data) {
return ProfileTaskLog.builder()
.id(data.getId())
.taskId(data.getTagValue(ProfileTaskLogRecord.TASK_ID))
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileThreadSnapshotQueryDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileThreadSnapshotQueryDAO.java
index fe4b815122..3072d6337f 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileThreadSnapshotQueryDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBProfileThreadSnapshotQueryDAO.java
@@ -19,6 +19,7 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import com.google.common.collect.ImmutableSet;
+import org.apache.skywalking.banyandb.v1.client.Element;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
@@ -120,7 +121,7 @@ public class BanyanDBProfileThreadSnapshotQueryDAO extends
AbstractBanyanDBDAO i
});
List<BasicTrace> basicTraces = new ArrayList<>();
- for (final RowEntity row : segmentRecordResp.getElements()) {
+ for (final Element row : segmentRecordResp.getElements()) {
BasicTrace basicTrace = new BasicTrace();
basicTrace.setSegmentId(row.getId());
diff --git
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBTraceQueryDAO.java
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBTraceQueryDAO.java
index f6491cd43d..641204bc3d 100644
---
a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBTraceQueryDAO.java
+++
b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBTraceQueryDAO.java
@@ -20,6 +20,7 @@ package
org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import com.google.common.collect.ImmutableSet;
import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
+import org.apache.skywalking.banyandb.v1.client.Element;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
@@ -152,7 +153,7 @@ public class BanyanDBTraceQueryDAO extends
AbstractBanyanDBDAO implements ITrace
return traceBrief;
}
- for (final RowEntity row : resp.getElements()) {
+ for (final Element row : resp.getElements()) {
BasicTrace basicTrace = new BasicTrace();
basicTrace.setSegmentId(row.getId());
diff --git a/skywalking-ui b/skywalking-ui
index 163de5e5cf..1278454148 160000
--- a/skywalking-ui
+++ b/skywalking-ui
@@ -1 +1 @@
-Subproject commit 163de5e5cf2a16a3e25e42d0aaa4b6b1c551a85f
+Subproject commit 127845414890fc3868dc9af73972f1c89e776366
diff --git a/test/e2e-v2/java-test-service/e2e-protocol/src/main/proto
b/test/e2e-v2/java-test-service/e2e-protocol/src/main/proto
index c2c9bbbff4..b8d5a5c27c 160000
--- a/test/e2e-v2/java-test-service/e2e-protocol/src/main/proto
+++ b/test/e2e-v2/java-test-service/e2e-protocol/src/main/proto
@@ -1 +1 @@
-Subproject commit c2c9bbbff43f0ad9f917a0a55285538a6b45739e
+Subproject commit b8d5a5c27c271303ff1d3911fb85a554352f4f23
diff --git a/test/e2e-v2/script/env b/test/e2e-v2/script/env
index e843c82ac6..c95c6a8327 100644
--- a/test/e2e-v2/script/env
+++ b/test/e2e-v2/script/env
@@ -23,6 +23,6 @@
SW_AGENT_CLIENT_JS_COMMIT=af0565a67d382b683c1dbd94c379b7080db61449
SW_AGENT_CLIENT_JS_TEST_COMMIT=4f1eb1dcdbde3ec4a38534bf01dded4ab5d2f016
SW_KUBERNETES_COMMIT_SHA=b670c41d94a82ddefcf466d54bab5c492d88d772
SW_ROVER_COMMIT=8550199e98c9f5a4b2058878a0a899ffb73fe461
-SW_BANYANDB_COMMIT=e7b08bea242e76c68950509529339995ac0646df
+SW_BANYANDB_COMMIT=1c19243df23f8350ea5c1542fd914c49e8698960
SW_CTL_COMMIT=0883266bfaa36612927b69e35781b64ea181758d