This is an automated email from the ASF dual-hosted git repository.
hui pushed a commit to branch lmh/TypeProviderOpt
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/lmh/TypeProviderOpt by this
push:
new 3d4e4c6396 fix CI
3d4e4c6396 is described below
commit 3d4e4c6396e13e65db1bfc0627c9f175f518bc60
Author: liuminghui233 <[email protected]>
AuthorDate: Wed Sep 7 23:35:51 2022 +0800
fix CI
---
.../main/java/org/apache/iotdb/db/mpp/common/NodeRef.java | 6 ++++++
.../apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java | 6 +++++-
.../iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java | 15 +++++++++++----
.../iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java | 4 +++-
.../distribution/SimpleFragmentParallelPlanner.java | 5 ++++-
5 files changed, 29 insertions(+), 7 deletions(-)
diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/common/NodeRef.java
b/server/src/main/java/org/apache/iotdb/db/mpp/common/NodeRef.java
index 761984d3dd..779e0eff91 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/common/NodeRef.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/common/NodeRef.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.db.mpp.common;
import org.apache.iotdb.db.mpp.plan.statement.StatementNode;
+import static java.lang.String.format;
import static java.lang.System.identityHashCode;
import static java.util.Objects.requireNonNull;
@@ -52,4 +53,9 @@ public class NodeRef<T extends StatementNode> {
public int hashCode() {
return identityHashCode(node);
}
+
+ @Override
+ public String toString() {
+ return format("@%s: %s", Integer.toHexString(identityHashCode(node)),
node);
+ }
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java
index a3f6cee2de..84c9988d0b 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java
@@ -497,6 +497,7 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
// generate result set header according to output expressions
DatasetHeader datasetHeader = analyzeOutput(analysis, outputExpressions);
+ analysis.setOutputExpressions(outputExpressions);
analysis.setRespDatasetHeader(datasetHeader);
// fetch partition information
@@ -658,9 +659,12 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
for (String deviceName :
deviceToTransformExpressionOfOneMeasurement.keySet()) {
Expression transformExpression =
deviceToTransformExpressionOfOneMeasurement.get(deviceName);
+ Expression transformExpressionWithoutAlias =
+
ExpressionAnalyzer.removeAliasFromExpression(transformExpression);
+ analyzeExpression(analysis, transformExpressionWithoutAlias);
deviceToTransformExpressions
.computeIfAbsent(deviceName, key -> new LinkedHashSet<>())
-
.add(ExpressionAnalyzer.removeAliasFromExpression(transformExpression));
+ .add(transformExpressionWithoutAlias);
deviceToMeasurementsMap
.computeIfAbsent(deviceName, key -> new LinkedHashSet<>())
.add(measurementExpressionWithoutAlias.getExpressionString());
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 fd221a0b43..0d04be5ff3 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
@@ -100,6 +100,7 @@ import java.util.stream.Collectors;
import static com.google.common.base.Preconditions.checkArgument;
import static
org.apache.iotdb.commons.conf.IoTDBConstant.MULTI_LEVEL_PATH_WILDCARD;
+import static
org.apache.iotdb.db.mpp.common.header.ColumnHeaderConstant.COLUMN_DEVICE;
public class LogicalPlanBuilder {
@@ -203,6 +204,7 @@ public class LogicalPlanBuilder {
Filter timeFilter,
GroupByTimeParameter groupByTimeParameter,
Set<Expression> aggregationExpressions,
+ Set<Expression> aggregationTransformExpressions,
Map<Expression, Set<Expression>> groupByLevelExpressions) {
boolean needCheckAscending = groupByTimeParameter == null;
Map<PartialPath, List<AggregationDescriptor>> ascendingAggregations = new
HashMap<>();
@@ -225,6 +227,7 @@ public class LogicalPlanBuilder {
timeFilter,
groupByTimeParameter);
updateTypeProvider(sourceExpressions);
+ updateTypeProvider(aggregationTransformExpressions);
return convergeAggregationSource(
sourceNodeList,
@@ -241,8 +244,9 @@ public class LogicalPlanBuilder {
Ordering scanOrder,
Filter timeFilter,
GroupByTimeParameter groupByTimeParameter,
- Set<Expression> aggregationExpressions,
List<Integer> measurementIndexes,
+ Set<Expression> aggregationExpressions,
+ Set<Expression> aggregationTransformExpressions,
Map<Expression, Set<Expression>> groupByLevelExpressions) {
checkArgument(
sourceExpressions.size() == measurementIndexes.size(),
@@ -275,6 +279,7 @@ public class LogicalPlanBuilder {
timeFilter,
groupByTimeParameter);
updateTypeProvider(sourceExpressions);
+ updateTypeProvider(aggregationTransformExpressions);
if (!curStep.isOutputPartial()) {
// update measurementIndexes
@@ -441,10 +446,12 @@ public class LogicalPlanBuilder {
List<Pair<Expression, String>> outputExpressions,
Map<String, List<Integer>> deviceToMeasurementIndexesMap,
Ordering mergeOrder) {
- List<String> outputColumnNames =
+ List<String> outputColumnNames = new ArrayList<>();
+ outputColumnNames.add(COLUMN_DEVICE);
+ outputColumnNames.addAll(
outputExpressions.stream()
.map(pair -> pair.getLeft().toString())
- .collect(Collectors.toList());
+ .collect(Collectors.toList()));
DeviceViewNode deviceViewNode =
new DeviceViewNode(
context.getQueryId().genPlanNodeId(),
@@ -460,7 +467,7 @@ public class LogicalPlanBuilder {
deviceViewNode.addChildDeviceNode(deviceName, subPlan);
}
- context.getTypeProvider().setType(ColumnHeaderConstant.COLUMN_DEVICE,
TSDataType.TEXT);
+ context.getTypeProvider().setType(COLUMN_DEVICE, TSDataType.TEXT);
updateTypeProvider(outputExpressions.stream().map(Pair::getLeft).collect(Collectors.toList()));
this.root = deviceViewNode;
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
index d1d5d46e20..b59e490ebe 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
@@ -301,8 +301,9 @@ public class LogicalPlanVisitor extends
StatementVisitor<PlanNode, MPPQueryConte
queryStatement.getResultTimeOrder(),
analysis.getGlobalTimeFilter(),
analysis.getGroupByTimeParameter(),
- aggregationExpressions,
measurementIndexes,
+ aggregationExpressions,
+ aggregationTransformExpressions,
analysis.getGroupByLevelExpressions());
if (queryStatement.isGroupByLevel()) {
planBuilder = // plan Having with GroupByLevel
@@ -330,6 +331,7 @@ public class LogicalPlanVisitor extends
StatementVisitor<PlanNode, MPPQueryConte
analysis.getGlobalTimeFilter(),
analysis.getGroupByTimeParameter(),
aggregationExpressions,
+ aggregationTransformExpressions,
analysis.getGroupByLevelExpressions());
if (queryStatement.isGroupByLevel()) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
index 2559f46e4b..e532877f6c 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
@@ -34,6 +34,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId;
import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeUtil;
import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.ExchangeNode;
import org.apache.iotdb.db.mpp.plan.planner.plan.node.sink.FragmentSinkNode;
+import org.apache.iotdb.db.mpp.plan.statement.crud.QueryStatement;
import org.apache.iotdb.tsfile.read.filter.basic.Filter;
import org.slf4j.Logger;
@@ -112,7 +113,9 @@ public class SimpleFragmentParallelPlanner implements
IFragmentParallelPlaner {
fragmentInstance.setDataRegionAndHost(regionReplicaSet);
fragmentInstance.setHostDataNode(selectTargetDataNode(regionReplicaSet));
-
fragmentInstance.getFragment().generateTypeProvider(queryContext.getTypeProvider());
+ if (analysis.getStatement() instanceof QueryStatement) {
+
fragmentInstance.getFragment().generateTypeProvider(queryContext.getTypeProvider());
+ }
instanceMap.putIfAbsent(fragment.getId(), fragmentInstance);
fragmentInstanceList.add(fragmentInstance);
}