This is an automated email from the ASF dual-hosted git repository.
xingtanzjr pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 1a18def1fc [IOTDB-3065]show lastest timeseries (#6063)
1a18def1fc is described below
commit 1a18def1fc641b451830363e6b030bf5394e6fd3
Author: ZhangHongYin <[email protected]>
AuthorDate: Tue May 31 14:37:18 2022 +0800
[IOTDB-3065]show lastest timeseries (#6063)
---
.../operator/schema/SchemaQueryMergeOperator.java | 2 -
.../schema/SchemaQueryOrderByHeatOperator.java | 161 +++++++++++++++++++++
.../apache/iotdb/db/mpp/plan/analyze/Analyzer.java | 33 ++++-
.../db/mpp/plan/analyze/ClusterSchemaFetcher.java | 1 -
.../db/mpp/plan/planner/LocalExecutionPlanner.java | 17 +++
.../db/mpp/plan/planner/LogicalPlanBuilder.java | 10 ++
.../iotdb/db/mpp/plan/planner/LogicalPlanner.java | 31 ++--
.../planner/distribution/ExchangeNodeAdder.java | 7 +
.../mpp/plan/planner/plan/node/PlanNodeType.java | 6 +-
.../db/mpp/plan/planner/plan/node/PlanVisitor.java | 5 +
.../node/metedata/read/SchemaQueryMergeNode.java | 4 +
.../metedata/read/SchemaQueryOrderByHeatNode.java | 99 +++++++++++++
.../metedata/read/TimeSeriesSchemaScanNode.java | 7 +
.../iotdb/db/mpp/plan/plan/LogicalPlannerTest.java | 6 +-
14 files changed, 373 insertions(+), 16 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryMergeOperator.java
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryMergeOperator.java
index 061e306e81..3517537321 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryMergeOperator.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryMergeOperator.java
@@ -50,8 +50,6 @@ public class SchemaQueryMergeOperator implements
ProcessOperator {
@Override
public TsBlock next() {
- // ToDo @xinzhongtianxia consider SHOW LATEST
-
for (int i = 0; i < children.size(); i++) {
if (!noMoreTsBlocks[i]) {
TsBlock tsBlock = children.get(i).next();
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryOrderByHeatOperator.java
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryOrderByHeatOperator.java
new file mode 100644
index 0000000000..24d3efd37d
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryOrderByHeatOperator.java
@@ -0,0 +1,161 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.mpp.execution.operator.schema;
+
+import org.apache.iotdb.db.mpp.common.header.HeaderConstant;
+import org.apache.iotdb.db.mpp.execution.operator.Operator;
+import org.apache.iotdb.db.mpp.execution.operator.OperatorContext;
+import org.apache.iotdb.db.mpp.execution.operator.process.ProcessOperator;
+import org.apache.iotdb.tsfile.read.common.block.TsBlock;
+import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder;
+import org.apache.iotdb.tsfile.utils.Binary;
+
+import com.google.common.util.concurrent.ListenableFuture;
+
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static java.util.Objects.requireNonNull;
+
+public class SchemaQueryOrderByHeatOperator implements ProcessOperator {
+
+ private final OperatorContext operatorContext;
+ private final Operator left;
+ private final Operator right;
+ private boolean isFinished = false;
+ private final List<TsBlock> leftResult;
+ private final List<TsBlock> rightResult;
+
+ public SchemaQueryOrderByHeatOperator(
+ OperatorContext operatorContext, Operator left, Operator right) {
+ this.operatorContext = requireNonNull(operatorContext, "operatorContext is
null");
+ this.left = requireNonNull(left, "left child operator is null");
+ this.right = requireNonNull(right, "right child operator is null");
+ this.leftResult = new ArrayList<>();
+ this.rightResult = new ArrayList<>();
+ }
+
+ @Override
+ public TsBlock next() {
+ isFinished = true;
+
+ TsBlockBuilder tsBlockBuilder =
+ new
TsBlockBuilder(HeaderConstant.showTimeSeriesHeader.getRespDataTypes());
+
+ // Step 1: get last point result
+ Map<String, Long> timeseriesToLastTimestamp = new HashMap<>();
+ for (TsBlock tsBlock : rightResult) {
+ if (null == tsBlock || tsBlock.isEmpty()) {
+ continue;
+ }
+ for (int i = 0; i < tsBlock.getPositionCount(); i++) {
+ String timeseries = tsBlock.getColumn(0).getBinary(i).toString();
+ long time = tsBlock.getTimeByIndex(i);
+ timeseriesToLastTimestamp.put(timeseries, time);
+ }
+ }
+
+ // Step 2: get last point timestamp to timeseries map
+ Map<Long, List<Object[]>> lastTimestampToTsSchema = new HashMap<>();
+ for (TsBlock tsBlock : leftResult) {
+ if (null == tsBlock || tsBlock.isEmpty()) {
+ continue;
+ }
+ TsBlock.TsBlockRowIterator tsBlockRowIterator =
tsBlock.getTsBlockRowIterator();
+ while (tsBlockRowIterator.hasNext()) {
+ Object[] line = tsBlockRowIterator.next();
+ String timeseries = line[0].toString();
+ long time = timeseriesToLastTimestamp.getOrDefault(timeseries, 0L);
+ if (!lastTimestampToTsSchema.containsKey(time)) {
+ lastTimestampToTsSchema.put(time, new ArrayList<>());
+ }
+ lastTimestampToTsSchema.get(time).add(line);
+ }
+ }
+
+ // Step 3: sort by last point's timestamp
+ List<Long> timestamps = new ArrayList<>(lastTimestampToTsSchema.keySet());
+ timestamps.sort(Comparator.reverseOrder());
+
+ // Step 4: generate result
+ for (Long time : timestamps) {
+ List<Object[]> rows = lastTimestampToTsSchema.get(time);
+ for (Object[] row : rows) {
+ tsBlockBuilder.getTimeColumnBuilder().writeLong(0L);
+ for (int i = 0; i <
HeaderConstant.showTimeSeriesHeader.getRespDataTypes().size(); i++) {
+ Object value = row[i];
+ if (null == value) {
+ tsBlockBuilder.getColumnBuilder(i).appendNull();
+ } else {
+ tsBlockBuilder.getColumnBuilder(i).writeBinary(new
Binary(value.toString()));
+ }
+ }
+ tsBlockBuilder.declarePosition();
+ }
+ }
+
+ return tsBlockBuilder.build();
+ }
+
+ @Override
+ public OperatorContext getOperatorContext() {
+ return operatorContext;
+ }
+
+ @Override
+ public ListenableFuture<Void> isBlocked() {
+ ListenableFuture<Void> blocked = left.isBlocked();
+ while (left.hasNext() && blocked.isDone()) {
+ leftResult.add(left.next());
+ blocked = left.isBlocked();
+ }
+ if (!blocked.isDone()) {
+ return blocked;
+ }
+ blocked = right.isBlocked();
+ while (right.hasNext() && blocked.isDone()) {
+ rightResult.add(right.next());
+ blocked = right.isBlocked();
+ }
+ if (!blocked.isDone()) {
+ return blocked;
+ }
+ return NOT_BLOCKED;
+ }
+
+ @Override
+ public boolean hasNext() {
+ return !isFinished;
+ }
+
+ @Override
+ public void close() throws Exception {
+ left.close();
+ right.close();
+ }
+
+ @Override
+ public boolean isFinished() {
+ return isFinished;
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java
index 985b6ec362..49d5389ed1 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java
@@ -866,7 +866,7 @@ public class Analyzer {
analysis.setStatement(statement);
SchemaTree schemaTree = new SchemaTree();
- schemaTree.setStorageGroups(schemaTree.getStorageGroups());
+ schemaTree.setStorageGroups(statement.getStorageGroups());
return analyzeLast(analysis, statement.getSelectedPaths(), schemaTree);
}
@@ -1136,6 +1136,37 @@ public class Analyzer {
partitionFetcher.getSchemaPartition(
new PathPatternTree(showTimeSeriesStatement.getPathPattern()));
analysis.setSchemaPartitionInfo(schemaPartitionInfo);
+
+ if (showTimeSeriesStatement.isOrderByHeat()) {
+ PathPatternTree patternTree = new
PathPatternTree(showTimeSeriesStatement.getPathPattern());
+ patternTree.constructTree();
+ // request schema fetch API
+ logger.info("{} fetch query schema...", getLogHeader());
+ SchemaTree schemaTree = schemaFetcher.fetchSchema(patternTree);
+ logger.info("{} fetch schema done", getLogHeader());
+ List<MeasurementPath> allSelectedPath = schemaTree.getAllMeasurement();
+
+ Set<Expression> sourceExpressions =
+ allSelectedPath.stream()
+ .map(TimeSeriesOperand::new)
+ .collect(Collectors.toCollection(LinkedHashSet::new));
+ analysis.setSourceExpressions(sourceExpressions);
+
+ Set<String> deviceSet =
+
allSelectedPath.stream().map(MeasurementPath::getDevice).collect(Collectors.toSet());
+ Map<String, List<DataPartitionQueryParam>> sgNameToQueryParamsMap =
new HashMap<>();
+ for (String devicePath : deviceSet) {
+ DataPartitionQueryParam queryParam = new DataPartitionQueryParam();
+ queryParam.setDevicePath(devicePath);
+ sgNameToQueryParamsMap
+ .computeIfAbsent(
+ schemaTree.getBelongedStorageGroup(devicePath), key -> new
ArrayList<>())
+ .add(queryParam);
+ }
+ DataPartition dataPartition =
partitionFetcher.getDataPartition(sgNameToQueryParamsMap);
+ analysis.setDataPartitionInfo(dataPartition);
+ }
+
analysis.setRespDatasetHeader(HeaderConstant.showTimeSeriesHeader);
return analysis;
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
index 874e46ac0e..a3867c174a 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
@@ -96,7 +96,6 @@ public class ClusterSchemaFetcher implements ISchemaFetcher {
}
private SchemaTree executeSchemaFetchQuery(SchemaFetchStatement
schemaFetchStatement) {
-
long queryId = SessionManager.getInstance().requestQueryId(false);
ExecutionResult executionResult =
coordinator.execute(schemaFetchStatement, queryId, null, "",
partitionFetcher, this);
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
index 73f89cf160..96e9dfdc1c 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
@@ -93,6 +93,7 @@ import
org.apache.iotdb.db.mpp.execution.operator.schema.NodePathsSchemaScanOper
import
org.apache.iotdb.db.mpp.execution.operator.schema.SchemaFetchMergeOperator;
import
org.apache.iotdb.db.mpp.execution.operator.schema.SchemaFetchScanOperator;
import
org.apache.iotdb.db.mpp.execution.operator.schema.SchemaQueryMergeOperator;
+import
org.apache.iotdb.db.mpp.execution.operator.schema.SchemaQueryOrderByHeatOperator;
import
org.apache.iotdb.db.mpp.execution.operator.schema.TimeSeriesCountOperator;
import
org.apache.iotdb.db.mpp.execution.operator.schema.TimeSeriesSchemaScanOperator;
import
org.apache.iotdb.db.mpp.execution.operator.source.AlignedSeriesAggregationScanOperator;
@@ -118,6 +119,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.NodePathsSch
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchMergeNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesCountNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesSchemaScanNode;
@@ -344,6 +346,21 @@ public class LocalExecutionPlanner {
return seriesAggregationScanOperator;
}
+ @Override
+ public Operator visitSchemaQueryOrderByHeat(
+ SchemaQueryOrderByHeatNode node, LocalExecutionPlanContext context) {
+ Operator left = node.getLeft().accept(this, context);
+ Operator right = node.getRight().accept(this, context);
+
+ OperatorContext operatorContext =
+ context.instanceContext.addOperatorContext(
+ context.getNextOperatorId(),
+ node.getPlanNodeId(),
+ SchemaQueryOrderByHeatOperator.class.getSimpleName());
+
+ return new SchemaQueryOrderByHeatOperator(operatorContext, left, right);
+ }
+
@Override
public Operator visitSchemaQueryScan(
SchemaQueryScanNode node, LocalExecutionPlanContext context) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java
index 2ad0ff99f8..d789e58892 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java
@@ -43,6 +43,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.NodePathsSch
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchMergeNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesCountNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesSchemaScanNode;
import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.AggregationNode;
@@ -586,6 +587,15 @@ public class LogicalPlanBuilder {
return this;
}
+ public LogicalPlanBuilder planSchemaQueryOrderByHeat(PlanNode lastPlanNode) {
+ SchemaQueryOrderByHeatNode node =
+ new SchemaQueryOrderByHeatNode(context.getQueryId().genPlanNodeId());
+ node.addChild(this.getRoot());
+ node.addChild(lastPlanNode);
+ this.root = node;
+ return this;
+ }
+
public LogicalPlanBuilder planSchemaFetchMerge() {
this.root = new SchemaFetchMergeNode(context.getQueryId().genPlanNodeId());
return this;
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java
index d681aba8f3..b6fa90d6c2 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java
@@ -425,17 +425,28 @@ public class LogicalPlanner {
public PlanNode visitShowTimeSeries(
ShowTimeSeriesStatement showTimeSeriesStatement, MPPQueryContext
context) {
LogicalPlanBuilder planBuilder = new LogicalPlanBuilder(context);
+ planBuilder =
+ planBuilder
+ .planTimeSeriesSchemaSource(
+ showTimeSeriesStatement.getPathPattern(),
+ showTimeSeriesStatement.getKey(),
+ showTimeSeriesStatement.getValue(),
+ showTimeSeriesStatement.getLimit(),
+ showTimeSeriesStatement.getOffset(),
+ showTimeSeriesStatement.isOrderByHeat(),
+ showTimeSeriesStatement.isContains(),
+ showTimeSeriesStatement.isPrefixPath())
+ .planSchemaQueryMerge(showTimeSeriesStatement.isOrderByHeat());
+ // show latest timeseries
+ if (showTimeSeriesStatement.isOrderByHeat()
+ && 0 !=
analysis.getDataPartitionInfo().getDataPartitionMap().size()) {
+ PlanNode lastPlanNode =
+ new LogicalPlanBuilder(context)
+ .planLast(analysis.getSourceExpressions(),
analysis.getGlobalTimeFilter())
+ .getRoot();
+ planBuilder = planBuilder.planSchemaQueryOrderByHeat(lastPlanNode);
+ }
return planBuilder
- .planTimeSeriesSchemaSource(
- showTimeSeriesStatement.getPathPattern(),
- showTimeSeriesStatement.getKey(),
- showTimeSeriesStatement.getValue(),
- showTimeSeriesStatement.getLimit(),
- showTimeSeriesStatement.getOffset(),
- showTimeSeriesStatement.isOrderByHeat(),
- showTimeSeriesStatement.isContains(),
- showTimeSeriesStatement.isPrefixPath())
- .planSchemaQueryMerge(showTimeSeriesStatement.isOrderByHeat())
.planOffset(showTimeSeriesStatement.getOffset())
.planLimit(showTimeSeriesStatement.getLimit())
.getRoot();
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
index 3d5582229a..4de9558101 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java
@@ -28,6 +28,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.CountSchemaM
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchMergeNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryScanNode;
import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.AggregationNode;
import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceMergeNode;
@@ -198,6 +199,12 @@ public class ExchangeNodeAdder extends
PlanVisitor<PlanNode, NodeGroupContext> {
return processMultiChildNode(node, context);
}
+ @Override
+ public PlanNode visitSchemaQueryOrderByHeat(
+ SchemaQueryOrderByHeatNode node, NodeGroupContext context) {
+ return processMultiChildNode(node, context);
+ }
+
@Override
public PlanNode visitGroupByLevel(GroupByLevelNode node, NodeGroupContext
context) {
return processMultiChildNode(node, context);
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
index 4abac32e03..20f404fbb2 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
@@ -30,6 +30,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.NodePathsSch
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchMergeNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesCountNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesSchemaScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.AlterTimeSeriesNode;
@@ -121,7 +122,8 @@ public enum PlanNodeType {
LAST_QUERY_SCAN((short) 46),
ALIGNED_LAST_QUERY_SCAN((short) 47),
LAST_QUERY_MERGE((short) 48),
- NODE_PATHS_COUNT((short) 49);
+ NODE_PATHS_COUNT((short) 49),
+ SCHEMA_QUERY_ORDER_BY_HEAT((short) 50);
private final short nodeType;
@@ -245,6 +247,8 @@ public enum PlanNodeType {
return LastQueryMergeNode.deserialize(buffer);
case 49:
return NodePathsCountNode.deserialize(buffer);
+ case 50:
+ return SchemaQueryOrderByHeatNode.deserialize(buffer);
default:
throw new IllegalArgumentException("Invalid node type: " + nodeType);
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
index 20ea3ab8aa..816dc19d8e 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
@@ -29,6 +29,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.NodePathsSch
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchMergeNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaFetchScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesCountNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesSchemaScanNode;
@@ -155,6 +156,10 @@ public abstract class PlanVisitor<R, C> {
return visitPlan(node, context);
}
+ public R visitSchemaQueryOrderByHeat(SchemaQueryOrderByHeatNode node, C
context) {
+ return visitPlan(node, context);
+ }
+
public R visitTimeSeriesSchemaScan(TimeSeriesSchemaScanNode node, C context)
{
return visitPlan(node, context);
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/SchemaQueryMergeNode.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/SchemaQueryMergeNode.java
index 5e4b8a043e..53422d2e78 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/SchemaQueryMergeNode.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/SchemaQueryMergeNode.java
@@ -38,6 +38,10 @@ public class SchemaQueryMergeNode extends
AbstractSchemaMergeNode {
this.orderByHeat = orderByHeat;
}
+ public boolean isOrderByHeat() {
+ return orderByHeat;
+ }
+
@Override
public PlanNode clone() {
return new SchemaQueryMergeNode(getPlanNodeId(), this.orderByHeat);
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/SchemaQueryOrderByHeatNode.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/SchemaQueryOrderByHeatNode.java
new file mode 100644
index 0000000000..a0cd47b48c
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/SchemaQueryOrderByHeatNode.java
@@ -0,0 +1,99 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read;
+
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeType;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.MultiChildNode;
+
+import com.google.common.collect.ImmutableList;
+
+import java.nio.ByteBuffer;
+import java.util.List;
+
+public class SchemaQueryOrderByHeatNode extends MultiChildNode {
+
+ public SchemaQueryOrderByHeatNode(PlanNodeId id) {
+ super(id);
+ }
+
+ /** show timeseries */
+ private PlanNode left;
+
+ /** last point */
+ private PlanNode right;
+
+ public PlanNode getLeft() {
+ return left;
+ }
+
+ public PlanNode getRight() {
+ return right;
+ }
+
+ @Override
+ public List<PlanNode> getChildren() {
+ return ImmutableList.of(left, right);
+ }
+
+ @Override
+ public void addChild(PlanNode child) {
+ if (child instanceof SchemaQueryMergeNode) {
+ left = child;
+ } else {
+ right = child;
+ }
+ }
+
+ @Override
+ public PlanNode clone() {
+ return new SchemaQueryOrderByHeatNode(getPlanNodeId());
+ }
+
+ @Override
+ public int allowedChildCount() {
+ return CHILD_COUNT_NO_LIMIT;
+ }
+
+ @Override
+ public List<String> getOutputColumnNames() {
+ return left.getOutputColumnNames();
+ }
+
+ @Override
+ public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
+ return visitor.visitSchemaQueryOrderByHeat(this, context);
+ }
+
+ @Override
+ protected void serializeAttributes(ByteBuffer byteBuffer) {
+ PlanNodeType.SCHEMA_QUERY_ORDER_BY_HEAT.serialize(byteBuffer);
+ }
+
+ public static SchemaQueryOrderByHeatNode deserialize(ByteBuffer byteBuffer) {
+ PlanNodeId planNodeId = PlanNodeId.deserialize(byteBuffer);
+ return new SchemaQueryOrderByHeatNode(planNodeId);
+ }
+
+ public String toString() {
+ return String.format("SchemaQueryOrderByHeatNode-%s", getPlanNodeId());
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/TimeSeriesSchemaScanNode.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/TimeSeriesSchemaScanNode.java
index 567b86c630..5f7ef53576 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/TimeSeriesSchemaScanNode.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/TimeSeriesSchemaScanNode.java
@@ -140,4 +140,11 @@ public class TimeSeriesSchemaScanNode extends
SchemaQueryScanNode {
public int hashCode() {
return Objects.hash(super.hashCode(), key, value, isContains, orderByHeat);
}
+
+ @Override
+ public String toString() {
+ return String.format(
+ "TimeSeriesSchemaScanNode-%s:[DataRegion: %s]",
+ this.getPlanNodeId(), this.getRegionReplicaSet());
+ }
}
diff --git
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/LogicalPlannerTest.java
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/LogicalPlannerTest.java
index 1482ca1f93..aa3ba7d7cf 100644
---
a/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/LogicalPlannerTest.java
+++
b/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/LogicalPlannerTest.java
@@ -38,6 +38,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.NodePathsCon
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.NodePathsCountNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.NodePathsSchemaScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.TimeSeriesSchemaScanNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.AlterTimeSeriesNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateAlignedTimeSeriesNode;
@@ -475,7 +476,10 @@ public class LogicalPlannerTest {
try {
LimitNode limitNode = (LimitNode) parseSQLToPlanNode(sql);
OffsetNode offsetNode = (OffsetNode) limitNode.getChild();
- SchemaQueryMergeNode metaMergeNode = (SchemaQueryMergeNode)
offsetNode.getChild();
+ SchemaQueryOrderByHeatNode schemaQueryOrderByHeatNode =
+ (SchemaQueryOrderByHeatNode) offsetNode.getChild();
+ SchemaQueryMergeNode metaMergeNode =
+ (SchemaQueryMergeNode)
schemaQueryOrderByHeatNode.getChildren().get(0);
metaMergeNode.getChildren().forEach(n ->
System.out.println(n.toString()));
TimeSeriesSchemaScanNode showTimeSeriesNode =
(TimeSeriesSchemaScanNode) metaMergeNode.getChildren().get(0);