This is an automated email from the ASF dual-hosted git repository.
xiangweiwei pushed a commit to branch aggregationVector
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/aggregationVector by this push:
new 732dd7a part2
732dd7a is described below
commit 732dd7a0c4b046cb381aeafaa0d41a1f721100ae
Author: Alima777 <[email protected]>
AuthorDate: Wed Sep 15 09:31:11 2021 +0800
part2
---
.../db/query/executor/AggregationExecutor.java | 130 +++++++++++++--------
.../query/executor/ISeriesAggregationExecutor.java | 9 ++
2 files changed, 91 insertions(+), 48 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/query/executor/AggregationExecutor.java
b/server/src/main/java/org/apache/iotdb/db/query/executor/AggregationExecutor.java
index c2e1a49..d2d45ae 100644
---
a/server/src/main/java/org/apache/iotdb/db/query/executor/AggregationExecutor.java
+++
b/server/src/main/java/org/apache/iotdb/db/query/executor/AggregationExecutor.java
@@ -19,7 +19,6 @@
package org.apache.iotdb.db.query.executor;
-import java.util.LinkedHashMap;
import org.apache.iotdb.db.conf.IoTDBConstant;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.engine.StorageEngine;
@@ -51,6 +50,7 @@ import org.apache.iotdb.tsfile.read.common.BatchData;
import org.apache.iotdb.tsfile.read.common.RowRecord;
import org.apache.iotdb.tsfile.read.expression.IExpression;
import org.apache.iotdb.tsfile.read.expression.impl.GlobalTimeExpression;
+import org.apache.iotdb.tsfile.read.filter.TimeFilter;
import org.apache.iotdb.tsfile.read.filter.basic.Filter;
import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet;
import org.apache.iotdb.tsfile.read.query.timegenerator.TimeGenerator;
@@ -59,6 +59,7 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
@@ -74,6 +75,7 @@ public class AggregationExecutor {
protected List<String> aggregations;
protected IExpression expression;
protected boolean ascending;
+ protected AggregateResult[] aggregateResultList;
/** aggregation batch calculation size. */
private int aggregateFetchSize;
@@ -85,6 +87,7 @@ public class AggregationExecutor {
this.expression = aggregationPlan.getExpression();
this.aggregateFetchSize =
IoTDBDescriptor.getInstance().getConfig().getBatchSize();
this.ascending = aggregationPlan.isAscending();
+ this.aggregateResultList = new AggregateResult[selectedSeries.size()];
}
/**
@@ -94,7 +97,7 @@ public class AggregationExecutor {
*/
public QueryDataSet executeWithoutValueFilter(
QueryContext context, AggregationPlan aggregationPlan)
- throws StorageEngineException, IOException, QueryProcessException {
+ throws StorageEngineException, QueryProcessException {
Filter timeFilter = null;
if (expression != null) {
@@ -102,19 +105,19 @@ public class AggregationExecutor {
}
// TODO use multi-thread
- List<ISeriesAggregationExecutor> seriesAggregationExecutor =
groupSeriesToExecutor(selectedSeries);
+ List<ISeriesAggregationExecutor> seriesAggregationExecutor =
+ groupSeriesToExecutor(selectedSeries);
AggregateResult[] aggregateResultList = new
AggregateResult[selectedSeries.size()];
// TODO-Cluster: group the paths by storage group to reduce communications
List<StorageGroupProcessor> list =
- StorageEngine.getInstance().mergeLock(new
ArrayList<>(pathToAggrIndexesMap.keySet()));
+ StorageEngine.getInstance().mergeLock(new ArrayList<>(selectedSeries));
try {
- for (Map.Entry<PartialPath, List<Integer>> entry :
pathToAggrIndexesMap.entrySet()) {
- aggregateOneSeries(
- entry,
- aggregateResultList,
-
aggregationPlan.getAllMeasurementsInDevice(entry.getKey().getDevice()),
- timeFilter,
- context);
+ for (ISeriesAggregationExecutor executor : seriesAggregationExecutor) {
+ List<AggregateResult> aggregateResults =
executor.getAggregateResult(context, timeFilter);
+ List<Integer> indexes = executor.getSeriesIndex();
+ for (int i = 0; i < indexes.size(); i++) {
+ aggregateResultList[indexes.get(i)] = aggregateResults.get(i);
+ }
}
} finally {
StorageEngine.getInstance().mergeUnLock(list);
@@ -131,30 +134,13 @@ public class AggregationExecutor {
* @param context query context
*/
protected void aggregateOneSeries(
- Map.Entry<PartialPath, List<Integer>> pathToAggrIndexes,
+ ISeriesAggregationExecutor executor,
AggregateResult[] aggregateResultList,
Set<String> allMeasurementsInDevice,
Filter timeFilter,
QueryContext context)
throws IOException, QueryProcessException, StorageEngineException {
- List<AggregateResult> ascAggregateResultList = new ArrayList<>();
- List<AggregateResult> descAggregateResultList = new ArrayList<>();
- boolean[] isAsc = new boolean[aggregateResultList.length];
- PartialPath seriesPath = pathToAggrIndexes.getKey();
- TSDataType tsDataType = dataTypes.get(pathToAggrIndexes.getValue().get(0));
-
- for (int i : pathToAggrIndexes.getValue()) {
- // construct AggregateResult
- AggregateResult aggregateResult =
- AggregateResultFactory.getAggrResultByName(aggregations.get(i),
tsDataType);
- if (aggregateResult.isAscending()) {
- ascAggregateResultList.add(aggregateResult);
- isAsc[i] = true;
- } else {
- descAggregateResultList.add(aggregateResult);
- }
- }
aggregateOneSeries(
seriesPath,
allMeasurementsInDevice,
@@ -514,30 +500,31 @@ public class AggregationExecutor {
}
/**
- * Merge same series, Group all the subSensors of one vector into one
VectorPartialPath and convert to ISeriesAggregationExecutor.
- * For example: Given: paths: s1, vector.s1, vector.s2, s1 and aggregations:
count, count, count, sum.
- * Then: SeriesAggregationExecutor s1 -> [0, 3],
VectorSeriesAggregationExecutor vector[s1, s2], Map{s1 -> 1, s2 -> 2}
+ * Merge same series, Group all the subSensors of one vector into one
VectorPartialPath and
+ * convert to ISeriesAggregationExecutor. For example: Given: paths: s1,
vector.s1, vector.s2, s1
+ * and aggregations: count, count, count, sum. Then:
SeriesAggregationExecutor s1 -> [0, 3],
+ * VectorSeriesAggregationExecutor vector[s1, s2], Map{s1 -> 1, s2 -> 2}
*
* @param selectedSeries selected series
* @return ISeriesAggregationExecutor list
*/
- private List<ISeriesAggregationExecutor> groupSeriesToExecutor(
- List<PartialPath> selectedSeries) {
+ private List<ISeriesAggregationExecutor>
groupSeriesToExecutor(List<PartialPath> selectedSeries) {
Map<String, ISeriesAggregationExecutor> pathToSeriesExecutor = new
LinkedHashMap<>();
for (int i = 0; i < selectedSeries.size(); i++) {
PartialPath path = selectedSeries.get(i);
String pathName = path.getFullPath();
+ ISeriesAggregationExecutor seriesExecutor;
if (!pathToSeriesExecutor.containsKey(pathName)) {
- ISeriesAggregationExecutor seriesExecutor =
getSeriesAggregationExecutor(path);
+ seriesExecutor = getSeriesAggregationExecutor(path);
pathToSeriesExecutor.put(pathName, seriesExecutor);
} else {
- // if pathToSeriesExecutor contains node
- ISeriesAggregationExecutor seriesExecutor =
pathToSeriesExecutor.get(pathName);
- if (seriesExecutor instanceof SeriesAggregationExecutor) {
- ((SeriesAggregationExecutor) seriesExecutor).addSeriesIndex(i);
- } else if (seriesExecutor instanceof VectorSeriesAggregationExecutor){
- ((VectorSeriesAggregationExecutor)
seriesExecutor).addSubSensorIndex((((VectorPartialPath) path).getSubSensor(0)),
i);
- }
+ seriesExecutor = pathToSeriesExecutor.get(pathName);
+ }
+ if (seriesExecutor instanceof SeriesAggregationExecutor) {
+ ((SeriesAggregationExecutor) seriesExecutor).addSeriesIndex(i);
+ } else if (seriesExecutor instanceof VectorSeriesAggregationExecutor) {
+ ((VectorSeriesAggregationExecutor) seriesExecutor)
+ .addSubSensorIndex((((VectorPartialPath) path).getSubSensor(0)),
i);
}
}
return new ArrayList<>(pathToSeriesExecutor.values());
@@ -552,10 +539,14 @@ public class AggregationExecutor {
}
public class SeriesAggregationExecutor implements ISeriesAggregationExecutor
{
+
private PartialPath path;
private TSDataType dataType;
private List<Integer> seriesIndex = new ArrayList<>();
- private IReaderByTimestamp readerByTimestamp;
+
+ private List<AggregateResult> ascAggregateResultList = new ArrayList<>();
+ private List<AggregateResult> descAggregateResultList = new ArrayList<>();
+ private boolean[] isAsc = new boolean[selectedSeries.size()];
public SeriesAggregationExecutor(PartialPath path) {
this.path = path.copy();
@@ -563,12 +554,40 @@ public class AggregationExecutor {
public void addSeriesIndex(int index) {
seriesIndex.add(index);
+ if (dataType == null) {
+ dataType = dataTypes.get(index);
+ }
+ AggregateResult aggregateResult =
+ AggregateResultFactory.getAggrResultByName(aggregations.get(index),
dataType);
+ if (aggregateResult.isAscending()) {
+ ascAggregateResultList.add(aggregateResult);
+ isAsc[index] = true;
+ } else {
+ descAggregateResultList.add(aggregateResult);
+ }
+ }
+
+ public List<Integer> getSeriesIndex() {
+ return seriesIndex;
+ }
+
+ @Override
+ public List<AggregateResult> getAggregateResult(QueryContext context,
TimeFilter timeFilter) {
+ // construct series reader without value filter
+ QueryDataSource queryDataSource =
+ QueryResourceManager.getInstance().getQueryDataSource(path, context,
timeFilter);
+ if (fileFilter != null) {
+ QueryUtils.filterQueryDataSource(queryDataSource, fileFilter);
+ }
+ // update filter by TTL
+ timeFilter = queryDataSource.updateFilterUsingTTL(timeFilter);
+ return null;
}
}
public class VectorSeriesAggregationExecutor implements
ISeriesAggregationExecutor {
private VectorPartialPath vectorPath;
- private List<TSDataType> dataTypes;
+ private List<TSDataType> subDataTypes;
private Map<String, List<Integer>> subSensorIndexMap = new HashMap<>();
private IReaderByTimestamp readerByTimestamp;
@@ -577,10 +596,25 @@ public class AggregationExecutor {
}
public void addSubSensorIndex(String subSensor, int index) {
- subSensorIndexMap.computeIfAbsent(subSensor, k -> {
- vectorPath.addSubSensor(k);
- return new ArrayList<>();
- }).add(index);
+ subSensorIndexMap
+ .computeIfAbsent(
+ subSensor,
+ k -> {
+ vectorPath.addSubSensor(k);
+ subDataTypes.add(dataTypes.get(index));
+ return new ArrayList<>();
+ })
+ .add(index);
+ }
+
+ @Override
+ public List<AggregateResult> getAggregateResult(QueryContext context,
TimeFilter timeFilter) {
+ return null;
+ }
+
+ @Override
+ public List<Integer> getSeriesIndex() {
+ return null;
}
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/query/executor/ISeriesAggregationExecutor.java
b/server/src/main/java/org/apache/iotdb/db/query/executor/ISeriesAggregationExecutor.java
index 120f4e7..4fca712 100644
---
a/server/src/main/java/org/apache/iotdb/db/query/executor/ISeriesAggregationExecutor.java
+++
b/server/src/main/java/org/apache/iotdb/db/query/executor/ISeriesAggregationExecutor.java
@@ -1,5 +1,14 @@
package org.apache.iotdb.db.query.executor;
+import org.apache.iotdb.db.query.aggregation.AggregateResult;
+import org.apache.iotdb.db.query.context.QueryContext;
+import org.apache.iotdb.tsfile.read.filter.TimeFilter;
+
+import java.util.List;
+
public interface ISeriesAggregationExecutor {
+ public List<AggregateResult> getAggregateResult(QueryContext context,
TimeFilter timeFilter);
+
+ public List<Integer> getSeriesIndex();
}