This is an automated email from the ASF dual-hosted git repository.
caogaofei pushed a commit to branch beyyes/agg_mergesort
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/beyyes/agg_mergesort by this
push:
new 2477c10d1f6 fix DeviceViewNode
2477c10d1f6 is described below
commit 2477c10d1f6f1fbc8e0ce9f14157dc9378da5c3c
Author: Beyyes <[email protected]>
AuthorDate: Wed Feb 7 10:46:47 2024 +0800
fix DeviceViewNode
---
...ator.java => AggregationMergeSortOperator.java} | 17 +++++++--------
.../operator/process/DeviceViewOperator.java | 3 ---
.../plan/planner/LogicalPlanBuilder.java | 10 ++++-----
.../plan/planner/OperatorTreeGenerator.java | 9 ++++----
.../planner/distribution/ExchangeNodeAdder.java | 24 +++++++++++++++++-----
.../plan/planner/distribution/SourceRewriter.java | 9 ++++----
.../plan/planner/plan/node/PlanGraphPrinter.java | 5 +++--
.../plan/planner/plan/node/PlanNodeType.java | 4 ++--
.../plan/planner/plan/node/PlanVisitor.java | 4 ++--
...SortNode.java => AggregationMergeSortNode.java} | 16 +++++++--------
.../planner/plan/node/process/DeviceViewNode.java | 8 ++++----
.../AlignByDeviceOrderByLimitOffsetTest.java | 5 +++--
12 files changed, 63 insertions(+), 51 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggMergeSortOperator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggregationMergeSortOperator.java
similarity index 93%
rename from
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggMergeSortOperator.java
rename to
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggregationMergeSortOperator.java
index 14348f90512..3efb8111624 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggMergeSortOperator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggregationMergeSortOperator.java
@@ -36,10 +36,11 @@ import com.google.common.util.concurrent.ListenableFuture;
import java.util.List;
/**
- * Since devices have been sorted by the merge order as expected, what {@link
AggMergeSortOperator}
- * need to do is traversing the device child operators, get all tsBlocks of
one device and transform
- * it to the form we need, adding the device column and allocating value
column to its expected
- * location, then get the next device operator until no next device.
+ * Since devices have been sorted by the merge order as expected, what {@link
+ * AggregationMergeSortOperator} need to do is traversing the device child
operators, get all
+ * tsBlocks of one device and transform it to the form we need, adding the
device column and
+ * allocating value column to its expected location, then get the next device
operator until no next
+ * device.
*
* <p>The deviceOperators can be aggregationSeriesScanOperator,
imeJoinOperator or
* seriesScanOperator that have not transformed the result form.
@@ -48,7 +49,7 @@ import java.util.List;
* [s1,s2,s3] is query, but only [s1, s3] exists in device1, then the column
of s2 will be filled
* with NullColumn.
*/
-public class AggMergeSortOperator implements ProcessOperator {
+public class AggregationMergeSortOperator implements ProcessOperator {
private final OperatorContext operatorContext;
// The size devices and deviceOperators should be the same.
@@ -63,7 +64,7 @@ public class AggMergeSortOperator implements ProcessOperator {
private int deviceIndex;
- public AggMergeSortOperator(
+ public AggregationMergeSortOperator(
OperatorContext operatorContext,
List<String> devices,
List<Operator> deviceOperators,
@@ -78,10 +79,6 @@ public class AggMergeSortOperator implements ProcessOperator
{
this.deviceIndex = 0;
}
- private String getCurDeviceName() {
- return devices.get(deviceIndex);
- }
-
private Operator getCurDeviceOperator() {
return deviceOperators.get(deviceIndex);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/DeviceViewOperator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/DeviceViewOperator.java
index 122b271096b..0dd6329ac40 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/DeviceViewOperator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/DeviceViewOperator.java
@@ -84,9 +84,6 @@ public class DeviceViewOperator implements ProcessOperator {
}
private Operator getCurDeviceOperator() {
- if (deviceIndex >= deviceOperators.size()) {
- System.out.printf("aaa");
- }
return deviceOperators.get(deviceIndex);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LogicalPlanBuilder.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LogicalPlanBuilder.java
index aee63275d50..3ecf47ad30a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LogicalPlanBuilder.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LogicalPlanBuilder.java
@@ -55,7 +55,7 @@ import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.Sche
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.TimeSeriesCountNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.TimeSeriesSchemaScanNode;
-import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggMergeSortNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationMergeSortNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.ColumnInjectNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceViewIntoNode;
@@ -460,8 +460,8 @@ public class LogicalPlanBuilder {
"Each aggregate should correspond to a column of output.");
boolean needCheckAscending = groupByTimeParameter == null;
- Map<PartialPath, List<AggregationDescriptor>> ascendingAggregations = new
HashMap<>();
- Map<PartialPath, List<AggregationDescriptor>> descendingAggregations = new
HashMap<>();
+ Map<PartialPath, List<AggregationDescriptor>> ascendingAggregations = new
LinkedHashMap<>();
+ Map<PartialPath, List<AggregationDescriptor>> descendingAggregations = new
LinkedHashMap<>();
Map<AggregationDescriptor, Integer> aggregationToIndexMap = new
HashMap<>();
Map<PartialPath, List<AggregationDescriptor>> countTimeAggregations = new
HashMap<>();
@@ -956,8 +956,8 @@ public class LogicalPlanBuilder {
Map<String, List<Integer>> deviceToMeasurementIndexesMap,
Map<String, PlanNode> deviceNameToSourceNodesMap,
long valueFilterLimit) {
- AggMergeSortNode aggMergeSortNode =
- new AggMergeSortNode(
+ AggregationMergeSortNode aggMergeSortNode =
+ new AggregationMergeSortNode(
context.getQueryId().genPlanNodeId(),
orderByParameter,
outputColumnNames,
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/OperatorTreeGenerator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/OperatorTreeGenerator.java
index 8661406f2a6..efe5996fc56 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/OperatorTreeGenerator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/OperatorTreeGenerator.java
@@ -45,7 +45,7 @@ import
org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceManage
import org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil;
import org.apache.iotdb.db.queryengine.execution.operator.Operator;
import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext;
-import
org.apache.iotdb.db.queryengine.execution.operator.process.AggMergeSortOperator;
+import
org.apache.iotdb.db.queryengine.execution.operator.process.AggregationMergeSortOperator;
import
org.apache.iotdb.db.queryengine.execution.operator.process.AggregationOperator;
import
org.apache.iotdb.db.queryengine.execution.operator.process.ColumnInjectOperator;
import
org.apache.iotdb.db.queryengine.execution.operator.process.DeviceViewIntoOperator;
@@ -167,7 +167,7 @@ import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.Sche
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaQueryScanNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.TimeSeriesCountNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.TimeSeriesSchemaScanNode;
-import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggMergeSortNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationMergeSortNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.ColumnInjectNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceViewIntoNode;
@@ -825,7 +825,8 @@ public class OperatorTreeGenerator extends
PlanVisitor<Operator, LocalExecutionP
}
@Override
- public Operator visitAggMergeSort(AggMergeSortNode node,
LocalExecutionPlanContext context) {
+ public Operator visitAggregationMergeSort(
+ AggregationMergeSortNode node, LocalExecutionPlanContext context) {
OperatorContext operatorContext =
context
.getDriverContext()
@@ -840,7 +841,7 @@ public class OperatorTreeGenerator extends
PlanVisitor<Operator, LocalExecutionP
.collect(Collectors.toList());
List<TSDataType> outputColumnTypes = getOutputColumnTypes(node,
context.getTypeProvider());
- return new AggMergeSortOperator(
+ return new AggregationMergeSortOperator(
operatorContext, node.getDevices(), children, deviceColumnIndex,
outputColumnTypes);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/ExchangeNodeAdder.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/ExchangeNodeAdder.java
index 03d02d58e8c..ced6ab74a0b 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/ExchangeNodeAdder.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/ExchangeNodeAdder.java
@@ -32,7 +32,7 @@ import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.Sche
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaQueryOrderByHeatNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaQueryScanNode;
-import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggMergeSortNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationMergeSortNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceMergeNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceViewNode;
@@ -202,7 +202,8 @@ public class ExchangeNodeAdder extends
PlanVisitor<PlanNode, NodeGroupContext> {
}
@Override
- public PlanNode visitAggMergeSort(AggMergeSortNode node, NodeGroupContext
context) {
+ public PlanNode visitAggregationMergeSort(
+ AggregationMergeSortNode node, NodeGroupContext context) {
return processMultiChildNode(node, context);
}
@@ -412,7 +413,7 @@ public class ExchangeNodeAdder extends
PlanVisitor<PlanNode, NodeGroupContext> {
return newNode;
}
- if (node instanceof AggMergeSortNode) {
+ if (node instanceof AggregationMergeSortNode) {
return processAggMergeSortNode(node, visitedChildren, context, newNode,
dataRegion);
}
@@ -512,7 +513,7 @@ public class ExchangeNodeAdder extends
PlanVisitor<PlanNode, NodeGroupContext> {
NodeGroupContext context,
MultiChildProcessNode newNode,
TRegionReplicaSet dataRegion) {
- AggMergeSortNode aggMergeSortNode = (AggMergeSortNode) node;
+ AggregationMergeSortNode aggMergeSortNode = (AggregationMergeSortNode)
node;
Map<TRegionReplicaSet, DeviceViewNode> regionTopKNodeMap = new HashMap<>();
for (PlanNode child : visitedChildren) {
TRegionReplicaSet region =
context.getNodeDistribution(child.getPlanNodeId()).region;
@@ -531,7 +532,7 @@ public class ExchangeNodeAdder extends
PlanVisitor<PlanNode, NodeGroupContext> {
new
NodeDistribution(NodeDistributionType.SAME_WITH_ALL_CHILDREN, region));
return childDeviceViewNode;
});
- String device = ((SeriesAggregationScanNode)
child).getSeriesPath().getDevice();
+ String device = getChildNodeDevice(child);
deviceViewNode
.getDeviceToMeasurementIndexesMap()
.put(device,
aggMergeSortNode.getDeviceToMeasurementIndexesMap().get(device));
@@ -556,6 +557,19 @@ public class ExchangeNodeAdder extends
PlanVisitor<PlanNode, NodeGroupContext> {
return newNode;
}
+ private String getChildNodeDevice(PlanNode child) {
+ String device;
+ if (child instanceof SeriesAggregationScanNode) {
+ device = ((SeriesAggregationScanNode) child).getSeriesPath().getDevice();
+ } else if (child instanceof AlignedSeriesAggregationScanNode) {
+ device = ((AlignedSeriesAggregationScanNode)
child).getAlignedPath().getDevice();
+ } else {
+ throw new UnsupportedOperationException(
+ String.format("Unsupported child node of AggMergeSortNode, node:
%s", child.getClass()));
+ }
+ return device;
+ }
+
@Override
public PlanNode visitSlidingWindowAggregation(
SlidingWindowAggregationNode node, NodeGroupContext context) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/SourceRewriter.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/SourceRewriter.java
index 3661d93b252..d152fcd4288 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/SourceRewriter.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/SourceRewriter.java
@@ -38,7 +38,7 @@ import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.Sche
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaFetchScanNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaQueryMergeNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metedata.read.SchemaQueryScanNode;
-import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggMergeSortNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationMergeSortNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceViewNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.GroupByLevelNode;
@@ -234,9 +234,10 @@ public class SourceRewriter extends
SimplePlanNodeRewriter<DistributionPlanConte
}
@Override
- public List<PlanNode> visitAggMergeSort(AggMergeSortNode node,
DistributionPlanContext context) {
- AggMergeSortNode newRoot =
- new AggMergeSortNode(
+ public List<PlanNode> visitAggregationMergeSort(
+ AggregationMergeSortNode node, DistributionPlanContext context) {
+ AggregationMergeSortNode newRoot =
+ new AggregationMergeSortNode(
context.queryContext.getQueryId().genPlanNodeId(),
node.getMergeOrderParameter(),
node.getOutputColumnNames(),
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanGraphPrinter.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanGraphPrinter.java
index 72ac2a5a581..6db8bb75fbd 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanGraphPrinter.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanGraphPrinter.java
@@ -23,7 +23,7 @@ import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
import org.apache.iotdb.commons.partition.DataPartition;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.db.queryengine.plan.expression.Expression;
-import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggMergeSortNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationMergeSortNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.ColumnInjectNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceMergeNode;
@@ -190,7 +190,8 @@ public class PlanGraphPrinter extends
PlanVisitor<List<String>, PlanGraphPrinter
}
@Override
- public List<String> visitAggMergeSort(AggMergeSortNode node, GraphContext
context) {
+ public List<String> visitAggregationMergeSort(
+ AggregationMergeSortNode node, GraphContext context) {
List<String> boxValue = new ArrayList<>();
boxValue.add(String.format("AggMergeSort-%s",
node.getPlanNodeId().getId()));
boxValue.add(String.format("AggDeviceCount: %d",
node.getDevices().size()));
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanNodeType.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanNodeType.java
index 1cf585eeb7d..3b40eb38a1d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanNodeType.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanNodeType.java
@@ -60,7 +60,7 @@ import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedC
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedDeleteDataNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedInsertNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedWriteSchemaNode;
-import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggMergeSortNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationMergeSortNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.ColumnInjectNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceMergeNode;
@@ -429,7 +429,7 @@ public enum PlanNodeType {
case 88:
return LeftOuterTimeJoinNode.deserialize(buffer);
case 89:
- return AggMergeSortNode.deserialize(buffer);
+ return AggregationMergeSortNode.deserialize(buffer);
default:
throw new IllegalArgumentException("Invalid node type: " + nodeType);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanVisitor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanVisitor.java
index 8b67a9a98bc..1091d9a352a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanVisitor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/PlanVisitor.java
@@ -57,7 +57,7 @@ import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedC
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedDeleteDataNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedInsertNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedWriteSchemaNode;
-import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggMergeSortNode;
+import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationMergeSortNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.AggregationNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.ColumnInjectNode;
import
org.apache.iotdb.db.queryengine.plan.planner.plan.node.process.DeviceMergeNode;
@@ -232,7 +232,7 @@ public abstract class PlanVisitor<R, C> {
return visitMultiChildProcess(node, context);
}
- public R visitAggMergeSort(AggMergeSortNode node, C context) {
+ public R visitAggregationMergeSort(AggregationMergeSortNode node, C context)
{
return visitMultiChildProcess(node, context);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/AggMergeSortNode.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/AggregationMergeSortNode.java
similarity index 94%
rename from
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/AggMergeSortNode.java
rename to
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/AggregationMergeSortNode.java
index 7ead4883973..58e1b9aed12 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/AggMergeSortNode.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/AggregationMergeSortNode.java
@@ -35,7 +35,7 @@ import java.util.List;
import java.util.Map;
import java.util.Objects;
-public class AggMergeSortNode extends MultiChildProcessNode {
+public class AggregationMergeSortNode extends MultiChildProcessNode {
// The result output order, which could sort by device and time.
// The size of this list is 2 and the first SortItem in this list has higher
priority.
@@ -51,7 +51,7 @@ public class AggMergeSortNode extends MultiChildProcessNode {
// not 0 because device is the first column
final Map<String, List<Integer>> deviceToMeasurementIndexesMap;
- public AggMergeSortNode(
+ public AggregationMergeSortNode(
PlanNodeId id,
OrderByParameter mergeOrderParameter,
List<String> outputColumnNames,
@@ -62,7 +62,7 @@ public class AggMergeSortNode extends MultiChildProcessNode {
this.deviceToMeasurementIndexesMap = deviceToMeasurementIndexesMap;
}
- public AggMergeSortNode(
+ public AggregationMergeSortNode(
PlanNodeId id,
OrderByParameter mergeOrderParameter,
List<String> outputColumnNames,
@@ -77,12 +77,12 @@ public class AggMergeSortNode extends MultiChildProcessNode
{
@Override
public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
- return visitor.visitAggMergeSort(this, context);
+ return visitor.visitAggregationMergeSort(this, context);
}
@Override
public PlanNode clone() {
- return new AggMergeSortNode(
+ return new AggregationMergeSortNode(
getPlanNodeId(),
mergeOrderParameter,
outputColumnNames,
@@ -106,7 +106,7 @@ public class AggMergeSortNode extends MultiChildProcessNode
{
if (!super.equals(o)) {
return false;
}
- AggMergeSortNode that = (AggMergeSortNode) o;
+ AggregationMergeSortNode that = (AggregationMergeSortNode) o;
return mergeOrderParameter.equals(that.mergeOrderParameter)
&& devices.equals(that.devices)
&& outputColumnNames.equals(that.outputColumnNames)
@@ -172,7 +172,7 @@ public class AggMergeSortNode extends MultiChildProcessNode
{
}
}
- public static AggMergeSortNode deserialize(ByteBuffer byteBuffer) {
+ public static AggregationMergeSortNode deserialize(ByteBuffer byteBuffer) {
OrderByParameter mergeOrderParameter =
OrderByParameter.deserialize(byteBuffer);
int columnSize = ReadWriteIOUtils.readInt(byteBuffer);
List<String> outputColumnNames = new ArrayList<>();
@@ -200,7 +200,7 @@ public class AggMergeSortNode extends MultiChildProcessNode
{
mapSize--;
}
PlanNodeId planNodeId = PlanNodeId.deserialize(byteBuffer);
- return new AggMergeSortNode(
+ return new AggregationMergeSortNode(
planNodeId, mergeOrderParameter, outputColumnNames, devices,
deviceToMeasurementIndexesMap);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/DeviceViewNode.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/DeviceViewNode.java
index 72fddc56896..4d4b9b67049 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/DeviceViewNode.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/process/DeviceViewNode.java
@@ -46,17 +46,17 @@ public class DeviceViewNode extends MultiChildProcessNode {
// The result output order, which could sort by device and time.
// The size of this list is 2 and the first SortItem in this list has higher
priority.
- protected final OrderByParameter mergeOrderParameter;
+ private final OrderByParameter mergeOrderParameter;
// The size devices and children should be the same.
- protected final List<String> devices = new ArrayList<>();
+ private final List<String> devices = new ArrayList<>();
// Device column and measurement columns in result output
- protected final List<String> outputColumnNames;
+ private final List<String> outputColumnNames;
// e.g. [s1,s2,s3] is query, but [s1, s3] exists in device1, then device1 ->
[1, 3], s1 is 1 but
// not 0 because device is the first column
- protected final Map<String, List<Integer>> deviceToMeasurementIndexesMap;
+ private final Map<String, List<Integer>> deviceToMeasurementIndexesMap;
public DeviceViewNode(
PlanNodeId id,
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/AlignByDeviceOrderByLimitOffsetTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/AlignByDeviceOrderByLimitOffsetTest.java
index bf9ac3e9861..1daecbab923 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/AlignByDeviceOrderByLimitOffsetTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/distribution/AlignByDeviceOrderByLimitOffsetTest.java
@@ -90,8 +90,7 @@ public class AlignByDeviceOrderByLimitOffsetTest {
@Test
public void orderByDeviceTest1() {
// no order by
- // sql = "select * from root.sg.d1, root.sg.d22 LIMIT 10 align by device";
- sql = "select first_value(s1) from root.sg.d1, root.sg.d22 align by
device";
+ sql = "select * from root.sg.d1, root.sg.d22 LIMIT 10 align by device";
analysis = Util.analyze(sql, context);
logicalPlanNode = Util.genLogicalPlan(analysis, context);
planner = new DistributionPlanner(analysis, new LogicalQueryPlan(context,
logicalPlanNode));
@@ -260,6 +259,8 @@ public class AlignByDeviceOrderByLimitOffsetTest {
*/
@Test
public void orderByDeviceTest4() {
+ // aggregation + order by device, no value filter
+
// aggregation + order by device, expression; with value filter
sql =
"select count(s1) from root.sg.d1, root.sg.d22 WHERE s2=1
having(count(s1)>1) LIMIT 5 align by device";