This is an automated email from the ASF dual-hosted git repository.
caogaofei pushed a commit to branch beyyes/multi_devices_fe
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/beyyes/multi_devices_fe by
this push:
new cfbc47b28d0 optimize for align by multi devices
cfbc47b28d0 is described below
commit cfbc47b28d0d19028977cc45f9c2603f8d74ec05
Author: Beyyes <[email protected]>
AuthorDate: Tue Oct 31 17:56:35 2023 +0800
optimize for align by multi devices
---
.../db/queryengine/plan/analyze/Analysis.java | 10 +
.../queryengine/plan/analyze/AnalyzeVisitor.java | 37 +-
.../plan/analyze/TemplatedDeviceAnalyze.java | 466 +++++++++++++++++----
3 files changed, 415 insertions(+), 98 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/Analysis.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/Analysis.java
index aee45a16193..0949abc857f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/Analysis.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/Analysis.java
@@ -80,6 +80,8 @@ public class Analysis {
// map from output column name (for every node) to its datatype
private final Map<NodeRef<Expression>, TSDataType> expressionTypes = new
LinkedHashMap<>();
+ private Template templateTypes;
+
private boolean finishQueryAfterAnalyze;
// potential fail status when finishQueryAfterAnalyze is true. If failStatus
is NULL, means no
@@ -354,6 +356,10 @@ public class Analysis {
if (expression.getExpressionType() == ExpressionType.NULL) {
return null;
}
+
+ // TODO add template optimization
+ // expression.get
+
TSDataType type = expressionTypes.get(NodeRef.of(expression));
checkArgument(type != null, "Expression is not analyzed: %s", expression);
return type;
@@ -676,6 +682,10 @@ public class Analysis {
return expressionTypes;
}
+ public Template getTemplateTypes() {
+ return this.templateTypes;
+ }
+
public void setOrderByExpressions(Set<Expression> orderByExpressions) {
this.orderByExpressions = orderByExpressions;
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java
index d40cd576319..5e8970d3a40 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/AnalyzeVisitor.java
@@ -262,7 +262,6 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
}
ISchemaTree schemaTree = analyzeSchema(queryStatement, analysis,
context);
-
logger.warn("----- Analyze analyzeSchema cost: {}ms",
System.currentTimeMillis() - startTime);
startTime = System.currentTimeMillis();
@@ -280,13 +279,16 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
List<Pair<Expression, String>> outputExpressions;
if (queryStatement.isAlignByDevice()) {
- // TODO
- // 并行 最上层比如加sort
- // fe 线程管理
- //
+
+ if (!queryStatement.isAggregationQuery()) {
+ analysis =
+ TemplatedDeviceAnalyze.visitQuery(analysis, queryStatement,
context, schemaTree);
+ // fetch partition information
+ analyzeDataPartition(analysis, queryStatement, schemaTree);
+ return analysis;
+ }
+
List<PartialPath> deviceList = analyzeFrom(queryStatement, schemaTree);
- logger.warn("----- Analyze analyzeFrom cost: {}ms",
System.currentTimeMillis() - startTime);
- startTime = System.currentTimeMillis();
if (canPushDownLimitOffsetInGroupByTimeForDevice(queryStatement)) {
// remove the device which won't appear in resultSet after
limit/offset
@@ -294,13 +296,10 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
}
analyzeDeviceToWhere(analysis, queryStatement, schemaTree, deviceList);
- logger.warn(
- "----- Analyze analyzeDeviceToWhere cost: {}ms",
- System.currentTimeMillis() - startTime);
- startTime = System.currentTimeMillis();
outputExpressions =
- TemplatedDeviceAnalyze.analyzeSelect(analysis, queryStatement,
schemaTree, deviceList);
+ TemplatedDeviceAnalyze.analyzeSelectUseTemplate(
+ analysis, queryStatement, schemaTree, deviceList);
logger.warn(
"----- Analyze analyzeSelect cost: {}ms",
System.currentTimeMillis() - startTime);
@@ -319,7 +318,7 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
analyzeDeviceToSourceTransform(analysis, queryStatement);
analyzeDeviceToSource(analysis, queryStatement);
logger.warn(
- "----- Analyze analyzeDeviceToSource cost: {}ms",
+ "----- Analyze analyzeDeviceToSource +
analyzeDeviceToSourceTransform cost: {}ms",
System.currentTimeMillis() - startTime);
startTime = System.currentTimeMillis();
@@ -663,14 +662,12 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
Set<PartialPath> deviceSet = new HashSet<>();
for (PartialPath devicePattern : devicePatternList) {
// get all matched devices
- // TODO isPrefixMatch可否设置为false? analyzeFrom能否直接返回schemaTree的全部devices?
deviceSet.addAll(
schemaTree.getMatchedDevices(devicePattern, false).stream()
.map(DeviceSchemaInfo::getDevicePath)
.collect(Collectors.toList()));
}
- // TODO 是否一定要排序? 最终的sourceNodeList已经会排序?
return queryStatement.getResultDeviceOrder() == Ordering.ASC
? deviceSet.stream().sorted().collect(Collectors.toList())
:
deviceSet.stream().sorted(Comparator.reverseOrder()).collect(Collectors.toList());
@@ -1348,7 +1345,7 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
}
}
- private static final String WHERE_WRONG_TYPE_ERROR_MSG =
+ static final String WHERE_WRONG_TYPE_ERROR_MSG =
"The output type of the expression in WHERE clause should be BOOLEAN,
actual data type: %s.";
private void analyzeDeviceToWhere(
@@ -1449,7 +1446,7 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
analyzeDeviceViewSpecialProcess(deviceViewOutputExpressions,
queryStatement, analysis));
}
- private boolean analyzeDeviceViewSpecialProcess(
+ static boolean analyzeDeviceViewSpecialProcess(
Set<Expression> deviceViewOutputExpressions,
QueryStatement queryStatement,
Analysis analysis) {
@@ -1509,7 +1506,7 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
analysis.setDeviceViewInputIndexesMap(deviceViewInputIndexesMap);
}
- private void checkDeviceViewInputUniqueness(Set<Expression>
outputExpressionsUnderDevice) {
+ static void checkDeviceViewInputUniqueness(Set<Expression>
outputExpressionsUnderDevice) {
Set<Expression> normalizedOutputExpressionsUnderDevice =
outputExpressionsUnderDevice.stream()
.map(ExpressionAnalyzer::normalizeExpression)
@@ -1521,7 +1518,7 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
}
}
- private void analyzeOutput(
+ static void analyzeOutput(
Analysis analysis,
QueryStatement queryStatement,
List<Pair<Expression, String>> outputExpressions) {
@@ -1655,7 +1652,7 @@ public class AnalyzeVisitor extends
StatementVisitor<Analysis, MPPQueryContext>
queryStatement.updateSortItems(orderByExpressions);
}
- private static TSDataType analyzeExpressionType(Analysis analysis,
Expression expression) {
+ static TSDataType analyzeExpressionType(Analysis analysis, Expression
expression) {
return analyzeExpression(analysis, expression);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/TemplatedDeviceAnalyze.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/TemplatedDeviceAnalyze.java
index 21c65eb2335..22cd2c20698 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/TemplatedDeviceAnalyze.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/TemplatedDeviceAnalyze.java
@@ -19,34 +19,56 @@
package org.apache.iotdb.db.queryengine.plan.analyze;
+import org.apache.iotdb.commons.path.MeasurementPath;
import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.db.exception.sql.MeasurementNotExistException;
+import org.apache.iotdb.db.exception.sql.SemanticException;
+import org.apache.iotdb.db.queryengine.common.MPPQueryContext;
+import org.apache.iotdb.db.queryengine.common.header.DatasetHeaderFactory;
+import org.apache.iotdb.db.queryengine.common.schematree.DeviceSchemaInfo;
import org.apache.iotdb.db.queryengine.common.schematree.ISchemaTree;
import org.apache.iotdb.db.queryengine.plan.expression.Expression;
import org.apache.iotdb.db.queryengine.plan.expression.leaf.TimeSeriesOperand;
-import org.apache.iotdb.db.queryengine.plan.statement.component.ResultColumn;
+import
org.apache.iotdb.db.queryengine.plan.expression.multi.FunctionExpression;
+import org.apache.iotdb.db.queryengine.plan.statement.component.Ordering;
import org.apache.iotdb.db.queryengine.plan.statement.crud.QueryStatement;
import org.apache.iotdb.db.schemaengine.template.ClusterTemplateManager;
import org.apache.iotdb.db.schemaengine.template.Template;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.utils.Pair;
+import org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.util.ArrayList;
+import java.util.Comparator;
import java.util.HashMap;
import java.util.HashSet;
+import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.stream.Collectors;
+import static com.google.common.base.Preconditions.checkState;
+import static
org.apache.iotdb.db.queryengine.common.header.ColumnHeaderConstant.ENDTIME;
import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.DEVICE_EXPRESSION;
import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.END_TIME_EXPRESSION;
-import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.analyzeAlias;
-import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.checkAliasUniqueness;
-import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.updateDeviceToSelectExpressions;
-import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.updateMeasurementToDeviceSelectExpressions;
-import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionAnalyzer.concatDeviceAndBindSchemaForExpression;
-import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionAnalyzer.toLowerCaseExpression;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.WHERE_WRONG_TYPE_ERROR_MSG;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.analyzeDeviceViewSpecialProcess;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.analyzeExpressionType;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.analyzeOutput;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeVisitor.checkDeviceViewInputUniqueness;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionAnalyzer.getMeasurementExpression;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionAnalyzer.normalizeExpression;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionAnalyzer.searchAggregationExpressions;
+import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionAnalyzer.searchSourceExpressions;
import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionTypeAnalyzer.analyzeExpression;
+import static
org.apache.iotdb.db.queryengine.plan.optimization.LimitOffsetPushDown.canPushDownLimitOffsetInGroupByTimeForDevice;
+import static
org.apache.iotdb.db.queryengine.plan.optimization.LimitOffsetPushDown.pushDownLimitOffsetInGroupByTimeForDevice;
/**
* This class provides accelerated implementation for multiple devices align
by device query. This
@@ -58,78 +80,192 @@ import static
org.apache.iotdb.db.queryengine.plan.analyze.ExpressionTypeAnalyze
*/
public class TemplatedDeviceAnalyze {
- protected static List<Pair<Expression, String>> analyzeSelect(
+ private static final Logger logger =
LoggerFactory.getLogger(TemplatedDeviceAnalyze.class);
+
+ /**
+ * 并行 最上层比如加sort fe 线程管理
+ *
+ * @param queryStatement
+ * @param context
+ * @return
+ */
+ static Analysis visitQuery(
Analysis analysis,
QueryStatement queryStatement,
- ISchemaTree schemaTree,
- List<PartialPath> deviceList) {
+ MPPQueryContext context,
+ ISchemaTree schemaTree) {
- Template template =
ClusterTemplateManager.getInstance().getAllTemplates().get(0);
+ long startTime = System.currentTimeMillis();
- List<Pair<Expression, String>> outputExpressions = new ArrayList<>();
- Map<String, Set<Expression>> deviceToSelectExpressions = new HashMap<>();
+ List<Pair<Expression, String>> outputExpressions;
- ColumnPaginationController paginationController =
- new ColumnPaginationController(
- queryStatement.getSeriesLimit(), queryStatement.getSeriesOffset(),
false);
-
- for (ResultColumn resultColumn :
queryStatement.getSelectComponent().getResultColumns()) {
- Expression selectExpression = resultColumn.getExpression();
-
- // select expression after removing wildcard, LinkedHashMap for
order-preserving
- Map<Expression, Map<String, Expression>>
measurementToDeviceSelectExpressions =
- new LinkedHashMap<>();
- for (PartialPath device : deviceList) {
- List<Expression> selectExpressionsOfOneDevice =
- concatDeviceAndBindSchemaForExpression(selectExpression, device,
schemaTree);
- if (selectExpressionsOfOneDevice.isEmpty()) {
- continue;
- }
+ List<PartialPath> deviceList = analyzeFrom(queryStatement, schemaTree);
+ logger.warn("----- Analyze analyzeFrom cost: {}ms",
System.currentTimeMillis() - startTime);
+ startTime = System.currentTimeMillis();
- updateMeasurementToDeviceSelectExpressions(
- analysis, measurementToDeviceSelectExpressions, device,
selectExpressionsOfOneDevice);
- }
+ if (canPushDownLimitOffsetInGroupByTimeForDevice(queryStatement)) {
+ // remove the device which won't appear in resultSet after limit/offset
+ deviceList = pushDownLimitOffsetInGroupByTimeForDevice(deviceList,
queryStatement);
+ }
- checkAliasUniqueness(resultColumn.getAlias(),
measurementToDeviceSelectExpressions);
+ analyzeDeviceToWhere(analysis, queryStatement, schemaTree, deviceList);
+ logger.warn(
+ "----- Analyze analyzeDeviceToWhere cost: {}ms",
System.currentTimeMillis() - startTime);
+ startTime = System.currentTimeMillis();
- for (Map.Entry<Expression, Map<String, Expression>> entry :
- measurementToDeviceSelectExpressions.entrySet()) {
- Expression measurementExpression = entry.getKey();
- Map<String, Expression> deviceToSelectExpressionsOfOneMeasurement =
entry.getValue();
+ outputExpressions =
+ TemplatedDeviceAnalyze.analyzeSelectUseTemplate(
+ analysis, queryStatement, schemaTree, deviceList);
- if (paginationController.hasCurOffset()) {
- paginationController.consumeOffset();
- } else if (paginationController.hasCurLimit()) {
- deviceToSelectExpressionsOfOneMeasurement
- .values()
- .forEach(expression -> analyzeExpression(analysis, expression));
+ logger.warn("----- Analyze analyzeSelect cost: {}ms",
System.currentTimeMillis() - startTime);
+ startTime = System.currentTimeMillis();
- // fix: devices used same template must have consistent type, no
need to
- // checkDataTypeConsistency
+ if (deviceList.isEmpty()) {
+ return finishQuery(queryStatement, analysis);
+ }
+ analysis.setDeviceList(deviceList);
- Expression lowerCaseMeasurementExpression =
toLowerCaseExpression(measurementExpression);
- analyzeExpression(analysis, lowerCaseMeasurementExpression);
+ // analyzeDeviceToGroupBy(analysis, queryStatement, schemaTree,
deviceList);
+ // analyzeDeviceToOrderBy(analysis, queryStatement, schemaTree,
deviceList);
+ // analyzeHaving(analysis, queryStatement, schemaTree, deviceList);
+ // analyzeDeviceToAggregation(analysis, queryStatement);
+ analyzeDeviceToSourceTransform(analysis, queryStatement);
+ analyzeDeviceToSource(analysis, queryStatement);
+ logger.warn(
+ "----- Analyze analyzeDeviceToSource + analyzeDeviceToSourceTransform
cost: {}ms",
+ System.currentTimeMillis() - startTime);
+ startTime = System.currentTimeMillis();
- outputExpressions.add(
- new Pair<>(
- lowerCaseMeasurementExpression,
- analyzeAlias(
- resultColumn.getAlias(),
- measurementExpression,
- lowerCaseMeasurementExpression,
- queryStatement)));
+ analyzeDeviceViewOutput(analysis, queryStatement);
+ analyzeDeviceViewInput(analysis, queryStatement);
- updateDeviceToSelectExpressions(
- analysis, deviceToSelectExpressions,
deviceToSelectExpressionsOfOneMeasurement);
+ logger.warn(
+ "----- Analyze analyzeDeviceView cost: {}ms",
System.currentTimeMillis() - startTime);
+ startTime = System.currentTimeMillis();
- paginationController.consumeLimit();
- } else {
- break;
- }
+ // analyzeInto(analysis, queryStatement, deviceList, outputExpressions);
+ // analyzeGroupByTime(analysis, queryStatement);
+ // analyzeFill(analysis, queryStatement);
+
+ // generate result set header according to output expressions
+ analyzeOutput(analysis, queryStatement, outputExpressions);
+
+ logger.warn(
+ "----- Analyze analyzeOutput+analyzeDataPartition cost: {}ms",
+ System.currentTimeMillis() - startTime);
+
+ return analysis;
+ }
+
+ private static List<PartialPath> analyzeFrom(
+ QueryStatement queryStatement, ISchemaTree schemaTree) {
+ // device path patterns in FROM clause
+ List<PartialPath> devicePatternList =
queryStatement.getFromComponent().getPrefixPaths();
+
+ Set<PartialPath> deviceSet = new HashSet<>();
+ for (PartialPath devicePattern : devicePatternList) {
+ // get all matched devices
+ // TODO isPrefixMatch可否设置为false? analyzeFrom能否直接返回schemaTree的全部devices?
+ deviceSet.addAll(
+ schemaTree.getMatchedDevices(devicePattern, false).stream()
+ .map(DeviceSchemaInfo::getDevicePath)
+ .collect(Collectors.toList()));
+ }
+
+ // TODO 是否一定要排序? 最终的sourceNodeList已经会排序?
+ return queryStatement.getResultDeviceOrder() == Ordering.ASC
+ ? deviceSet.stream().sorted().collect(Collectors.toList())
+ :
deviceSet.stream().sorted(Comparator.reverseOrder()).collect(Collectors.toList());
+ }
+
+ private static void analyzeDeviceToWhere(
+ Analysis analysis,
+ QueryStatement queryStatement,
+ ISchemaTree schemaTree,
+ List<PartialPath> deviceSet) {
+ if (!queryStatement.hasWhere()) {
+ return;
+ }
+
+ Map<String, Expression> deviceToWhereExpression = new HashMap<>();
+ Iterator<PartialPath> deviceIterator = deviceSet.iterator();
+ while (deviceIterator.hasNext()) {
+ PartialPath devicePath = deviceIterator.next();
+ Expression whereExpression;
+ try {
+ whereExpression =
+ normalizeExpression(analyzeWhereSplitByDevice(queryStatement,
devicePath, schemaTree));
+ } catch (MeasurementNotExistException e) {
+ logger.warn(
+ "Meets MeasurementNotExistException in analyzeDeviceToWhere when
executing align by device, "
+ + "error msg: {}",
+ e.getMessage());
+ deviceIterator.remove();
+ continue;
+ }
+
+ TSDataType outputType = analyzeExpressionType(analysis, whereExpression);
+ if (outputType != TSDataType.BOOLEAN) {
+ throw new SemanticException(String.format(WHERE_WRONG_TYPE_ERROR_MSG,
outputType));
}
+
+ deviceToWhereExpression.put(devicePath.getFullPath(), whereExpression);
+ }
+ analysis.setDeviceToWhereExpression(deviceToWhereExpression);
+ }
+
+ private static Expression analyzeWhereSplitByDevice(
+ QueryStatement queryStatement, PartialPath devicePath, ISchemaTree
schemaTree) {
+ List<Expression> conJunctions =
+ ExpressionAnalyzer.concatDeviceAndBindSchemaForPredicate(
+ queryStatement.getWhereCondition().getPredicate(), devicePath,
schemaTree, true);
+ return ExpressionUtils.constructQueryFilter(
+ conJunctions.stream().distinct().collect(Collectors.toList()));
+ }
+
+ private static Analysis finishQuery(QueryStatement queryStatement, Analysis
analysis) {
+ if (queryStatement.isSelectInto()) {
+ analysis.setRespDatasetHeader(
+
DatasetHeaderFactory.getSelectIntoHeader(queryStatement.isAlignByDevice()));
+ }
+ if (queryStatement.isLastQuery()) {
+ analysis.setRespDatasetHeader(DatasetHeaderFactory.getLastQueryHeader());
}
+ analysis.setFinishQueryAfterAnalyze(true);
+ return analysis;
+ }
- removeDevicesWithoutMeasurements(deviceList, deviceToSelectExpressions,
analysis);
+ static List<Pair<Expression, String>> analyzeSelectUseTemplate(
+ Analysis analysis,
+ QueryStatement queryStatement,
+ ISchemaTree schemaTree,
+ List<PartialPath> deviceList) {
+
+ Template template =
ClusterTemplateManager.getInstance().getAllTemplates().get(0);
+
+ List<Pair<Expression, String>> outputExpressions = new ArrayList<>();
+ Map<String, Set<Expression>> deviceToSelectExpressions = new HashMap<>();
+
+ for (Map.Entry<String, IMeasurementSchema> entry :
template.getSchemaMap().entrySet()) {
+ String measurementName = entry.getKey();
+ IMeasurementSchema measurementSchema = entry.getValue();
+ TimeSeriesOperand measurementPath =
+ new TimeSeriesOperand(
+ new MeasurementPath(new String[] {measurementName},
measurementSchema));
+ analyzeExpression(analysis, measurementPath);
+ outputExpressions.add(new Pair<>(measurementPath, null));
+ for (PartialPath devicePath : deviceList) {
+ // TODO how to determine whether a device is aligned device
+ TimeSeriesOperand fullPath =
+ new TimeSeriesOperand(
+ new MeasurementPath(
+ devicePath.concatNode(measurementName), measurementSchema,
true));
+ analyzeExpression(analysis, fullPath);
+ deviceToSelectExpressions
+ .computeIfAbsent(devicePath.getFullPath(), k -> new
LinkedHashSet<>())
+ .add(fullPath);
+ }
+ }
Set<Expression> selectExpressions = new LinkedHashSet<>();
selectExpressions.add(DEVICE_EXPRESSION);
@@ -140,28 +276,202 @@ public class TemplatedDeviceAnalyze {
analysis.setSelectExpressions(selectExpressions);
analysis.setDeviceToSelectExpressions(deviceToSelectExpressions);
-
return outputExpressions;
}
- private static void removeDevicesWithoutMeasurements(
- List<PartialPath> deviceList,
- Map<String, Set<Expression>> deviceToSelectExpressions,
- Analysis analysis) {
- // remove devices without measurements to compute
- Set<PartialPath> noMeasurementDevices = new HashSet<>();
- for (PartialPath device : deviceList) {
- if (!deviceToSelectExpressions.containsKey(device.getFullPath())) {
- noMeasurementDevices.add(device);
+ private static void analyzeDeviceToSourceTransform(
+ Analysis analysis, QueryStatement queryStatement) {
+ if (queryStatement.isAggregationQuery()) {
+ Map<String, Set<Expression>> deviceToSourceTransformExpressions =
+ analysis.getDeviceToSourceTransformExpressions();
+ Map<String, Set<Expression>> deviceToAggregationExpressions =
+ analysis.getDeviceToAggregationExpressions();
+
+ for (Map.Entry<String, Set<Expression>> entry :
deviceToAggregationExpressions.entrySet()) {
+ String deviceName = entry.getKey();
+ Set<Expression> aggregationExpressions = entry.getValue();
+
+ Set<Expression> sourceTransformExpressions =
+ deviceToSourceTransformExpressions.computeIfAbsent(
+ deviceName, k -> new LinkedHashSet<>());
+
+ for (Expression expression : aggregationExpressions) {
+ // if count_time aggregation exist, it can exist only one
count_time(*)
+ if (queryStatement.isCountTimeAggregation()) {
+ for (Expression countTimeSourceExpression :
+ ((FunctionExpression) expression).getCountTimeExpressions()) {
+
+ analyzeExpressionType(analysis, countTimeSourceExpression);
+ sourceTransformExpressions.add(countTimeSourceExpression);
+ }
+ } else {
+ // We just process first input Expression of AggregationFunction,
+ // keep other input Expressions as origin
+ // If AggregationFunction need more than one input series,
+ // we need to reconsider the process of it
+ sourceTransformExpressions.add(expression.getExpressions().get(0));
+ }
+ }
+
+ if (queryStatement.hasGroupByExpression()) {
+
sourceTransformExpressions.add(analysis.getDeviceToGroupByExpression().get(deviceName));
+ }
+ }
+ } else {
+ updateDeviceToSourceTransformAndOutputExpressions(
+ analysis, analysis.getDeviceToSelectExpressions());
+ if (queryStatement.hasOrderByExpression()) {
+ updateDeviceToSourceTransformAndOutputExpressions(
+ analysis, analysis.getDeviceToOrderByExpressions());
+ }
+ }
+ }
+
+ private static void updateDeviceToSourceTransformAndOutputExpressions(
+ Analysis analysis, Map<String, Set<Expression>>
deviceToSelectExpressions) {
+ // two maps to be updated
+ Map<String, Set<Expression>> deviceToSourceTransformExpressions =
+ analysis.getDeviceToSourceTransformExpressions();
+ Map<String, Set<Expression>> deviceToOutputExpressions =
+ analysis.getDeviceToOutputExpressions();
+
+ for (Map.Entry<String, Set<Expression>> entry :
deviceToSelectExpressions.entrySet()) {
+ String deviceName = entry.getKey();
+ Set<Expression> expressions = entry.getValue();
+
+ Set<Expression> normalizedExpressions = new LinkedHashSet<>();
+ for (Expression expression : expressions) {
+ Expression normalizedExpression = normalizeExpression(expression);
+ analyzeExpressionType(analysis, normalizedExpression);
+
+ normalizedExpressions.add(normalizedExpression);
+ }
+ deviceToOutputExpressions
+ .computeIfAbsent(deviceName, key -> new LinkedHashSet<>())
+ .addAll(expressions);
+ deviceToSourceTransformExpressions
+ .computeIfAbsent(deviceName, key -> new LinkedHashSet<>())
+ .addAll(normalizedExpressions);
+ }
+ }
+
+ private static void analyzeDeviceViewOutput(Analysis analysis,
QueryStatement queryStatement) {
+ Set<Expression> selectExpressions = analysis.getSelectExpressions();
+ Set<Expression> deviceViewOutputExpressions = new LinkedHashSet<>();
+ if (queryStatement.isAggregationQuery()) {
+ deviceViewOutputExpressions.add(DEVICE_EXPRESSION);
+ if (queryStatement.isOutputEndTime()) {
+ deviceViewOutputExpressions.add(END_TIME_EXPRESSION);
+ }
+ for (Expression selectExpression : selectExpressions) {
+
deviceViewOutputExpressions.addAll(searchAggregationExpressions(selectExpression));
+ }
+ if (queryStatement.hasHaving()) {
+ deviceViewOutputExpressions.addAll(
+ searchAggregationExpressions(analysis.getHavingExpression()));
+ }
+ if (queryStatement.hasOrderByExpression()) {
+ for (Expression orderByExpression : analysis.getOrderByExpressions()) {
+
deviceViewOutputExpressions.addAll(searchAggregationExpressions(orderByExpression));
+ }
+ }
+ } else {
+ deviceViewOutputExpressions.addAll(selectExpressions);
+ if (queryStatement.hasOrderByExpression()) {
+ deviceViewOutputExpressions.addAll(analysis.getOrderByExpressions());
+ }
+ }
+ analysis.setDeviceViewOutputExpressions(deviceViewOutputExpressions);
+ analysis.setDeviceViewSpecialProcess(
+ analyzeDeviceViewSpecialProcess(deviceViewOutputExpressions,
queryStatement, analysis));
+ }
+
+ private static void analyzeDeviceViewInput(Analysis analysis, QueryStatement
queryStatement) {
+ List<String> deviceViewOutputColumns =
+ analysis.getDeviceViewOutputExpressions().stream()
+ .map(Expression::getOutputSymbol)
+ .collect(Collectors.toList());
+
+ Map<String, Set<String>> deviceToOutputColumnsMap = new LinkedHashMap<>();
+ Map<String, Set<Expression>> deviceToOutputExpressions =
+ analysis.getDeviceToOutputExpressions();
+ for (Map.Entry<String, Set<Expression>> deviceOutputExpressionEntry :
+ deviceToOutputExpressions.entrySet()) {
+ Set<Expression> outputExpressionsUnderDevice =
deviceOutputExpressionEntry.getValue();
+ checkDeviceViewInputUniqueness(outputExpressionsUnderDevice);
+
+ Set<String> outputColumns = new LinkedHashSet<>();
+ if (queryStatement.isOutputEndTime()) {
+ outputColumns.add(ENDTIME);
+ }
+ for (Expression expression : outputExpressionsUnderDevice) {
+ outputColumns.add(getMeasurementExpression(expression,
analysis).getOutputSymbol());
+ }
+ deviceToOutputColumnsMap.put(deviceOutputExpressionEntry.getKey(),
outputColumns);
+ }
+
+ Map<String, List<Integer>> deviceViewInputIndexesMap = new HashMap<>();
+ for (Map.Entry<String, Set<String>> deviceOutputColumnsEntry :
+ deviceToOutputColumnsMap.entrySet()) {
+ String deviceName = deviceOutputColumnsEntry.getKey();
+ List<String> outputsUnderDevice = new
ArrayList<>(deviceOutputColumnsEntry.getValue());
+
+ List<Integer> indexes = new ArrayList<>();
+ for (String output : outputsUnderDevice) {
+ int index = deviceViewOutputColumns.indexOf(output);
+ checkState(
+ index >= 1, "output column '%s' is not stored in %s", output,
deviceViewOutputColumns);
+ indexes.add(index);
+ }
+ deviceViewInputIndexesMap.put(deviceName, indexes);
+ }
+ analysis.setDeviceViewInputIndexesMap(deviceViewInputIndexesMap);
+ }
+
+ private static void analyzeDeviceToSource(Analysis analysis, QueryStatement
queryStatement) {
+ Map<String, Set<Expression>> deviceToSourceExpressions = new HashMap<>();
+ Map<String, Set<Expression>> deviceToSourceTransformExpressions =
+ analysis.getDeviceToSourceTransformExpressions();
+
+ for (Map.Entry<String, Set<Expression>> entry :
deviceToSourceTransformExpressions.entrySet()) {
+ String deviceName = entry.getKey();
+ Set<Expression> sourceTransformExpressions = entry.getValue();
+
+ Set<Expression> sourceExpressions = new LinkedHashSet<>();
+ sourceTransformExpressions.forEach(
+ expression ->
sourceExpressions.addAll(searchSourceExpressions(expression)));
+
+ deviceToSourceExpressions.put(deviceName, sourceExpressions);
+ }
+
+ if (queryStatement.hasWhere()) {
+ Map<String, Expression> deviceToWhereExpression =
analysis.getDeviceToWhereExpression();
+ for (Map.Entry<String, Expression> deviceWhereExpressionEntry :
+ deviceToWhereExpression.entrySet()) {
+ String deviceName = deviceWhereExpressionEntry.getKey();
+ Expression whereExpression = deviceWhereExpressionEntry.getValue();
+ deviceToSourceExpressions
+ .computeIfAbsent(deviceName, key -> new LinkedHashSet<>())
+ .addAll(searchSourceExpressions(whereExpression));
}
}
- deviceList.removeAll(noMeasurementDevices);
- // when the select expression of any device is empty,
- // the where expression map also need remove this device
- if (analysis.getDeviceToWhereExpression() != null) {
- noMeasurementDevices.forEach(
- devicePath ->
analysis.getDeviceToWhereExpression().remove(devicePath.getFullPath()));
+ Map<String, List<String>> outputDeviceToQueriedDevicesMap = new
LinkedHashMap<>();
+ for (Map.Entry<String, Set<Expression>> entry :
deviceToSourceExpressions.entrySet()) {
+ String deviceName = entry.getKey();
+ Set<Expression> sourceExpressionsUnderDevice = entry.getValue();
+ Set<String> queriedDevices = new HashSet<>();
+ for (Expression expression : sourceExpressionsUnderDevice) {
+
queriedDevices.add(ExpressionAnalyzer.getDeviceNameInSourceExpression(expression));
+ }
+ if (queriedDevices.size() > 1) {
+ throw new SemanticException(
+ "Cross-device queries are not supported in ALIGN BY DEVICE
queries.");
+ }
+ outputDeviceToQueriedDevicesMap.put(deviceName, new
ArrayList<>(queriedDevices));
}
+
+ analysis.setDeviceToSourceExpressions(deviceToSourceExpressions);
+
analysis.setOutputDeviceToQueriedDevicesMap(outputDeviceToQueriedDevicesMap);
}
}