diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/input/Entity.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/input/Entity.java index a78116156..117e2d38b 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/input/Entity.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/input/Entity.java @@ -18,6 +18,8 @@ package org.apache.skywalking.oap.server.core.query.input; +import java.util.Objects; +import lombok.AccessLevel; import lombok.Getter; import lombok.Setter; import org.apache.skywalking.oap.server.core.analysis.IDManager; @@ -29,7 +31,7 @@ import org.apache.skywalking.oap.server.core.query.enumeration.Scope; * @since 8.0.0 */ @Setter -@Getter +@Getter(AccessLevel.PRIVATE) public class Entity { /** *
@@ -48,7 +50,7 @@ public class Entity {
      * Normal service is the service having installed agent or metrics reported directly. Unnormal service is
      * conjectural service, usually detected by the agent.
      */
-    private boolean normal;
+    private Boolean normal;
     private String serviceInstanceName;
     private String endpointName;
 
@@ -57,7 +59,7 @@ public class Entity {
      * Normal service is the service having installed agent or metrics reported directly. Unnormal service is
      * conjectural service, usually detected by the agent.
      */
-    private boolean destNormal;
+    private Boolean destNormal;
     private String destServiceInstanceName;
     private String destEndpointName;
 
@@ -65,6 +67,36 @@ public class Entity {
         return Scope.Service.equals(scope);
     }
 
+    /**
+     * @return true if the entity field is valid. The graphql definition couldn't provide the strict validation, because
+     * the required fields are according to the scope.
+     */
+    public boolean isValid() {
+        switch (scope) {
+            case All:
+                return true;
+            case Service:
+                return Objects.nonNull(serviceName) && Objects.nonNull(normal);
+            case ServiceInstance:
+                return Objects.nonNull(serviceName) && Objects.nonNull(serviceInstanceName) && Objects.nonNull(normal);
+            case Endpoint:
+                return Objects.nonNull(serviceName) && Objects.nonNull(endpointName) && Objects.nonNull(normal);
+            case ServiceRelation:
+                return Objects.nonNull(serviceName) && Objects.nonNull(destServiceName)
+                    && Objects.nonNull(normal) && Objects.nonNull(destNormal);
+            case ServiceInstanceRelation:
+                return Objects.nonNull(serviceName) && Objects.nonNull(destServiceName)
+                    && Objects.nonNull(serviceInstanceName) && Objects.nonNull(destServiceInstanceName)
+                    && Objects.nonNull(normal) && Objects.nonNull(destNormal);
+            case EndpointRelation:
+                return Objects.nonNull(serviceName) && Objects.nonNull(endpointName)
+                    && Objects.nonNull(serviceInstanceName) && Objects.nonNull(endpointName)
+                    && Objects.nonNull(normal) && Objects.nonNull(destNormal);
+            default:
+                return false;
+        }
+    }
+
     /**
      * @return entity id based on the definition.
      */
diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/Model.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/Model.java
index ec469ced3..0df4b0a5f 100644
--- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/Model.java
+++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/Model.java
@@ -21,14 +21,12 @@ package org.apache.skywalking.oap.server.core.storage.model;
 import java.util.List;
 import lombok.EqualsAndHashCode;
 import lombok.Getter;
-import lombok.RequiredArgsConstructor;
 import org.apache.skywalking.oap.server.core.analysis.DownSampling;
 
 /**
  * The model definition of a logic entity.
  */
 @Getter
-@RequiredArgsConstructor
 @EqualsAndHashCode
 public class Model {
     private final String name;
@@ -38,5 +36,22 @@ public class Model {
     private final DownSampling downsampling;
     private final boolean record;
     private final boolean superDataset;
+    private final boolean isTimeSeries;
 
+    public Model(final String name,
+                 final List columns,
+                 final List extraQueryIndices,
+                 final int scopeId,
+                 final DownSampling downsampling,
+                 final boolean record,
+                 final boolean superDataset) {
+        this.name = name;
+        this.columns = columns;
+        this.extraQueryIndices = extraQueryIndices;
+        this.scopeId = scopeId;
+        this.downsampling = downsampling;
+        this.isTimeSeries = !DownSampling.None.equals(downsampling);
+        this.record = record;
+        this.superDataset = superDataset;
+    }
 }
diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java
index 51742dd93..038d169d9 100644
--- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java
+++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/ttl/DataTTLKeeperTimer.java
@@ -86,6 +86,9 @@ public enum DataTTLKeeperTimer {
 
     private void execute(Model model) {
         try {
+            if (!model.isTimeSeries()) {
+                return;
+            }
             moduleManager.find(StorageModule.NAME)
                          .provider()
                          .getService(IHistoryDeleteDAO.class)
diff --git a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricQuery.java b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricQuery.java
index 393fde52a..109ac0693 100644
--- a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricQuery.java
+++ b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricQuery.java
@@ -100,9 +100,9 @@ public class MetricQuery implements GraphQLQueryResolver {
         final List metricsValues = query.readLabeledMetricsValues(condition, labels, duration);
         List response = new ArrayList<>(metricsValues.size());
         labels.forEach(l -> metricsValues.stream()
-            .filter(m -> m.getLabel().equals(l))
-            .findAny()
-            .ifPresent(values -> response.add(values.getValues())));
+                                         .filter(m -> m.getLabel().equals(l))
+                                         .findAny()
+                                         .ifPresent(values -> response.add(values.getValues())));
         return response;
     }
 
@@ -119,9 +119,9 @@ public class MetricQuery implements GraphQLQueryResolver {
         final List metricsValues = query.readLabeledMetricsValues(condition, labels, duration);
         List response = new ArrayList<>(metricsValues.size());
         labels.forEach(l -> metricsValues.stream()
-            .filter(m -> m.getLabel().equals(l))
-            .findAny()
-            .ifPresent(values -> response.add(values.getValues())));
+                                         .filter(m -> m.getLabel().equals(l))
+                                         .findAny()
+                                         .ifPresent(values -> response.add(values.getValues())));
         return response;
     }
 
@@ -160,6 +160,11 @@ public class MetricQuery implements GraphQLQueryResolver {
     private static class MockEntity extends Entity {
         private final String id;
 
+        @Override
+        public boolean isValid() {
+            return true;
+        }
+
         @Override
         public String buildId() {
             return id;
diff --git a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricsQuery.java b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricsQuery.java
index 2e257d13d..af7207724 100644
--- a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricsQuery.java
+++ b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/MetricsQuery.java
@@ -103,7 +103,7 @@ public class MetricsQuery implements GraphQLQueryResolver {
      * Read metrics single value in the duration of required metrics
      */
     public int readMetricsValue(MetricsCondition condition, Duration duration) throws IOException {
-        if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName()))) {
+        if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName())) || !condition.getEntity().isValid()) {
             return 0;
         }
         return getMetricsQueryService().readMetricsValue(condition, duration);
@@ -113,11 +113,13 @@ public class MetricsQuery implements GraphQLQueryResolver {
      * Read time-series values in the duration of required metrics
      */
     public MetricsValues readMetricsValues(MetricsCondition condition, Duration duration) throws IOException {
-        if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName()))) {
+        if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName())) || !condition.getEntity().isValid()) {
             final List pointOfTimes = duration.assembleDurationPoints();
             MetricsValues values = new MetricsValues();
             pointOfTimes.forEach(pointOfTime -> {
-                String id = pointOfTime.id(condition.getEntity().buildId());
+                String id = pointOfTime.id(
+                    condition.getEntity().isValid() ? condition.getEntity().buildId() : "ILLEGAL_ENTITY"
+                );
                 final KVInt kvInt = new KVInt();
                 kvInt.setId(id);
                 kvInt.setValue(0);
@@ -146,14 +148,16 @@ public class MetricsQuery implements GraphQLQueryResolver {
     public List readLabeledMetricsValues(MetricsCondition condition,
                                                         List labels,
                                                         Duration duration) throws IOException {
-        if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName()))) {
+        if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName())) || !condition.getEntity().isValid()) {
             final List pointOfTimes = duration.assembleDurationPoints();
 
             List labeledValues = new ArrayList<>(labels.size());
             labels.forEach(label -> {
                 MetricsValues values = new MetricsValues();
                 pointOfTimes.forEach(pointOfTime -> {
-                    String id = pointOfTime.id(condition.getEntity().buildId());
+                    String id = pointOfTime.id(
+                        condition.getEntity().isValid() ? condition.getEntity().buildId() : "ILLEGAL_ENTITY"
+                    );
                     final KVInt kvInt = new KVInt();
                     kvInt.setId(id);
                     kvInt.setValue(0);
@@ -180,14 +184,16 @@ public class MetricsQuery implements GraphQLQueryResolver {
      * 
*/ public HeatMap readHeatMap(MetricsCondition condition, Duration duration) throws IOException { - if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName()))) { + if (MetricsType.UNKNOWN.equals(typeOfMetrics(condition.getName())) || !condition.getEntity().isValid()) { DataTable emptyData = new DataTable(); emptyData.put("0", 0L); final String rawdata = emptyData.toStorageData(); final HeatMap heatMap = new HeatMap(); final List pointOfTimes = duration.assembleDurationPoints(); pointOfTimes.forEach(pointOfTime -> { - String id = pointOfTime.id(condition.getEntity().buildId()); + String id = pointOfTime.id( + condition.getEntity().isValid() ? condition.getEntity().buildId() : "ILLEGAL_ENTITY" + ); heatMap.buildColumn(id, rawdata, 0); }); return heatMap; diff --git a/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol b/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol index 0b38f4e16..4c1d1d996 160000 --- a/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol +++ b/oap-server/server-query-plugin/query-graphql-plugin/src/main/resources/query-protocol @@ -1 +1 @@ -Subproject commit 0b38f4e1620ef66e6d9121845b2a55b090d2020f +Subproject commit 4c1d1d996f6baece949fce90c676647b52e25620 diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java index ac69acfe0..2a131601e 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java @@ -19,13 +19,10 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base; import com.google.gson.Gson; - import java.io.IOException; import java.util.HashMap; import java.util.Map; - import lombok.extern.slf4j.Slf4j; - import org.apache.skywalking.apm.util.StringUtil; import org.apache.skywalking.oap.server.core.storage.StorageException; import org.apache.skywalking.oap.server.core.storage.model.Model; @@ -56,7 +53,10 @@ public class StorageEsInstaller extends ModelInstaller { protected boolean isExists(Model model) throws StorageException { ElasticSearchClient esClient = (ElasticSearchClient) client; try { - String timeSeriesIndexName = TimeSeriesUtils.latestWriteIndexName(model); + String timeSeriesIndexName = + model.isTimeSeries() ? + TimeSeriesUtils.latestWriteIndexName(model) : + model.getName(); return esClient.isExistsTemplate(model.getName()) && esClient.isExistsIndex(timeSeriesIndexName); } catch (IOException e) { throw new StorageException(e.getMessage()); @@ -73,22 +73,28 @@ public class StorageEsInstaller extends ModelInstaller { .toString()); try { - if (!esClient.isExistsTemplate(model.getName())) { - boolean isAcknowledged = esClient.createTemplate(model.getName(), settings, mapping); - log.info( - "create {} index template finished, isAcknowledged: {}", model.getName(), isAcknowledged); + String indexName; + if (!model.isTimeSeries()) { + indexName = model.getName(); + } else { + if (!esClient.isExistsTemplate(model.getName())) { + boolean isAcknowledged = esClient.createTemplate(model.getName(), settings, mapping); + log.info( + "create {} index template finished, isAcknowledged: {}", model.getName(), isAcknowledged); + if (!isAcknowledged) { + throw new StorageException("create " + model.getName() + " index template failure, "); + } + } + indexName = TimeSeriesUtils.latestWriteIndexName(model); + } + if (!esClient.isExistsIndex(indexName)) { + boolean isAcknowledged = esClient.createIndex(indexName); + log.info("create {} index finished, isAcknowledged: {}", indexName, isAcknowledged); if (!isAcknowledged) { - throw new StorageException("create " + model.getName() + " index template failure, "); - } - } - String timeSeriesIndexName = TimeSeriesUtils.latestWriteIndexName(model); - if (!esClient.isExistsIndex(timeSeriesIndexName)) { - boolean isAcknowledged = esClient.createIndex(timeSeriesIndexName); - log.info("create {} index finished, isAcknowledged: {}", timeSeriesIndexName, isAcknowledged); - if (!isAcknowledged) { - throw new StorageException("create " + timeSeriesIndexName + " time series index failure, "); + throw new StorageException("create " + indexName + " time series index failure, "); } } + } catch (IOException e) { throw new StorageException(e.getMessage()); } diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TopNRecordsQueryDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TopNRecordsQueryDAO.java index fe34e0ef3..82a665594 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TopNRecordsQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TopNRecordsQueryDAO.java @@ -49,12 +49,12 @@ public class H2TopNRecordsQueryDAO implements ITopNRecordsQueryDAO { List parameters = new ArrayList<>(10); if (StringUtil.isNotEmpty(condition.getParentService())) { - sql.append(" service_id = ? "); + sql.append(" service_id = ? and"); final String serviceId = IDManager.ServiceID.buildId(condition.getParentService(), condition.isNormal()); parameters.add(serviceId); } - sql.append(" and ").append(TopN.TIME_BUCKET).append(" >= ?"); + sql.append(" ").append(TopN.TIME_BUCKET).append(" >= ?"); parameters.add(duration.getStartTimeBucketInSec()); sql.append(" and ").append(TopN.TIME_BUCKET).append(" <= ?"); parameters.add(duration.getEndTimeBucketInSec());