This is an automated email from the ASF dual-hosted git repository.
jackietien pushed a commit to branch ty/AggPerf
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/ty/AggPerf by this push:
new f36dae8b483 opt
f36dae8b483 is described below
commit f36dae8b483941b98c632f1f03ee19a123e53550
Author: JackieTien97 <[email protected]>
AuthorDate: Mon Oct 28 18:28:46 2024 +0800
opt
---
.../execution/operator/source/FileLoaderUtils.java | 4 +---
.../TableAggregationTableScanOperator.java | 6 ++++-
.../source/relational/TableScanOperator.java | 12 +++++++---
.../plan/planner/TableOperatorGenerator.java | 18 ++++++++++++--
.../dataregion/read/QueryDataSource.java | 3 +++
.../apache/iotdb/commons/path/AlignedFullPath.java | 28 ++++++++++++++++++++++
6 files changed, 62 insertions(+), 9 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/FileLoaderUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/FileLoaderUtils.java
index e4f4459f41d..1941f3c739f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/FileLoaderUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/FileLoaderUtils.java
@@ -51,7 +51,6 @@ import org.apache.tsfile.read.reader.IPageReader;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
-import java.util.HashSet;
import java.util.List;
import java.util.Set;
@@ -264,8 +263,7 @@ public class FileLoaderUtils {
// the order of timeSeriesMetadata list is same as subSensorList's order
TimeSeriesMetadataCache cache = TimeSeriesMetadataCache.getInstance();
List<String> valueMeasurementList = alignedPath.getMeasurementList();
- Set<String> allSensors = new HashSet<>(valueMeasurementList);
- allSensors.add("");
+ Set<String> allSensors = alignedPath.getAllSensors();
boolean isDebug = context.isDebug();
String filePath = resource.getTsFilePath();
IDeviceID deviceId = alignedPath.getDeviceId();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableAggregationTableScanOperator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableAggregationTableScanOperator.java
index e5b3e2b1b3f..96678c0752a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableAggregationTableScanOperator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableAggregationTableScanOperator.java
@@ -59,6 +59,7 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
+import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@@ -92,6 +93,7 @@ public class TableAggregationTableScanOperator extends
AbstractSeriesAggregation
private final SeriesScanOptions seriesScanOptions;
private final List<String> measurementColumnNames;
+ private final Set<String> allSensors;
private final List<IMeasurementSchema> measurementSchemas;
@@ -120,6 +122,7 @@ public class TableAggregationTableScanOperator extends
AbstractSeriesAggregation
Ordering scanOrder,
SeriesScanOptions seriesScanOptions,
List<String> measurementColumnNames,
+ Set<String> allSensors,
List<IMeasurementSchema> measurementSchemas,
int maxTsBlockLineNum,
int measurementCount,
@@ -158,6 +161,7 @@ public class TableAggregationTableScanOperator extends
AbstractSeriesAggregation
this.scanOrder = scanOrder;
this.seriesScanOptions = seriesScanOptions;
this.measurementColumnNames = measurementColumnNames;
+ this.allSensors = allSensors;
this.measurementSchemas = measurementSchemas;
this.measurementColumnTSDataTypes =
measurementSchemas.stream().map(IMeasurementSchema::getType).collect(Collectors.toList());
@@ -266,7 +270,7 @@ public class TableAggregationTableScanOperator extends
AbstractSeriesAggregation
}
AlignedFullPath alignedPath =
- constructAlignedPath(deviceEntry, measurementColumnNames,
measurementSchemas);
+ constructAlignedPath(deviceEntry, measurementColumnNames,
measurementSchemas, allSensors);
this.seriesScanUtil =
new AlignedSeriesScanUtil(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableScanOperator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableScanOperator.java
index e282f1b3e63..dc04af63a96 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableScanOperator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/TableScanOperator.java
@@ -49,6 +49,7 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
+import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@@ -75,6 +76,8 @@ public class TableScanOperator extends
AbstractSeriesScanOperator {
private final List<String> measurementColumnNames;
+ private final Set<String> allSensors;
+
private final List<IMeasurementSchema> measurementSchemas;
private final List<TSDataType> measurementColumnTSDataTypes;
@@ -98,6 +101,7 @@ public class TableScanOperator extends
AbstractSeriesScanOperator {
Ordering scanOrder,
SeriesScanOptions seriesScanOptions,
List<String> measurementColumnNames,
+ Set<String> allSensors,
List<IMeasurementSchema> measurementSchemas,
int maxTsBlockLineNum) {
this.sourceId = sourceId;
@@ -109,6 +113,7 @@ public class TableScanOperator extends
AbstractSeriesScanOperator {
this.scanOrder = scanOrder;
this.seriesScanOptions = seriesScanOptions;
this.measurementColumnNames = measurementColumnNames;
+ this.allSensors = allSensors;
this.measurementSchemas = measurementSchemas;
this.measurementColumnTSDataTypes =
measurementSchemas.stream().map(IMeasurementSchema::getType).collect(Collectors.toList());
@@ -294,7 +299,7 @@ public class TableScanOperator extends
AbstractSeriesScanOperator {
private AlignedSeriesScanUtil constructAlignedSeriesScanUtil(DeviceEntry
deviceEntry) {
AlignedFullPath alignedPath =
- constructAlignedPath(deviceEntry, measurementColumnNames,
measurementSchemas);
+ constructAlignedPath(deviceEntry, measurementColumnNames,
measurementSchemas, allSensors);
return new AlignedSeriesScanUtil(
alignedPath,
@@ -308,9 +313,10 @@ public class TableScanOperator extends
AbstractSeriesScanOperator {
public static AlignedFullPath constructAlignedPath(
DeviceEntry deviceEntry,
List<String> measurementColumnNames,
- List<IMeasurementSchema> measurementSchemas) {
+ List<IMeasurementSchema> measurementSchemas,
+ Set<String> allSensors) {
return new AlignedFullPath(
- deviceEntry.getDeviceID(), measurementColumnNames, measurementSchemas);
+ deviceEntry.getDeviceID(), measurementColumnNames, measurementSchemas,
allSensors);
}
@Override
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TableOperatorGenerator.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TableOperatorGenerator.java
index d3150c2f731..1e4e4c7410d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TableOperatorGenerator.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/TableOperatorGenerator.java
@@ -382,6 +382,9 @@ public class TableOperatorGenerator extends
PlanVisitor<Operator, LocalExecution
context.getTypeProvider().getTemplatedInfo().getLimitValue(),
maxTsBlockLineNum);
}
+ Set<String> allSensors = new HashSet<>(measurementColumnNames);
+ // for time column
+ allSensors.add("");
TableScanOperator tableScanOperator =
new TableScanOperator(
operatorContext,
@@ -392,6 +395,7 @@ public class TableOperatorGenerator extends
PlanVisitor<Operator, LocalExecution
node.getScanOrder(),
scanOptionsBuilder.build(),
measurementColumnNames,
+ allSensors,
measurementSchemas,
maxTsBlockLineNum);
@@ -400,7 +404,10 @@ public class TableOperatorGenerator extends
PlanVisitor<Operator, LocalExecution
for (int i = 0, size = node.getDeviceEntries().size(); i < size; i++) {
AlignedFullPath alignedPath =
constructAlignedPath(
- node.getDeviceEntries().get(i), measurementColumnNames,
measurementSchemas);
+ node.getDeviceEntries().get(i),
+ measurementColumnNames,
+ measurementSchemas,
+ allSensors);
((DataDriverContext) context.getDriverContext()).addPath(alignedPath);
}
@@ -1632,6 +1639,9 @@ public class TableOperatorGenerator extends
PlanVisitor<Operator, LocalExecution
convertPredicateToFilter(pushDownPredicate, measurementColumnNames,
columnSchemaMap));
}
+ Set<String> allSensors = new HashSet<>(measurementColumnNames);
+ // for time column
+ allSensors.add("");
TableAggregationTableScanOperator aggTableScanOperator =
new TableAggregationTableScanOperator(
node.getPlanNodeId(),
@@ -1642,6 +1652,7 @@ public class TableOperatorGenerator extends
PlanVisitor<Operator, LocalExecution
scanAscending ? Ordering.ASC : Ordering.DESC,
scanOptionsBuilder.build(),
measurementColumnNames,
+ allSensors,
measurementSchemas,
TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(),
measurementColumnCount,
@@ -1659,7 +1670,10 @@ public class TableOperatorGenerator extends
PlanVisitor<Operator, LocalExecution
for (int i = 0, size = node.getDeviceEntries().size(); i < size; i++) {
AlignedFullPath alignedPath =
constructAlignedPath(
- node.getDeviceEntries().get(i), measurementColumnNames,
measurementSchemas);
+ node.getDeviceEntries().get(i),
+ measurementColumnNames,
+ measurementSchemas,
+ allSensors);
((DataDriverContext) context.getDriverContext()).addPath(alignedPath);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/read/QueryDataSource.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/read/QueryDataSource.java
index 4748cc20e58..704bdeb6902 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/read/QueryDataSource.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/read/QueryDataSource.java
@@ -200,6 +200,9 @@ public class QueryDataSource implements IQueryDataSource {
}
public void fillOrderIndexes(IDeviceID deviceId, boolean ascending) {
+ if (unseqResources == null || unseqResources.isEmpty()) {
+ return;
+ }
TreeMap<Long, List<Integer>> orderTimeToIndexMap =
ascending ? new TreeMap<>() : new TreeMap<>(descendingComparator);
int index = 0;
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/path/AlignedFullPath.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/path/AlignedFullPath.java
index 0ed64897438..581e448c253 100644
---
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/path/AlignedFullPath.java
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/path/AlignedFullPath.java
@@ -24,8 +24,12 @@ import org.apache.tsfile.file.metadata.IDeviceID;
import org.apache.tsfile.utils.RamUsageEstimator;
import org.apache.tsfile.write.schema.IMeasurementSchema;
+import javax.annotation.Nullable;
+
+import java.util.HashSet;
import java.util.List;
import java.util.Objects;
+import java.util.Set;
public class AlignedFullPath implements IFullPath {
@@ -38,12 +42,25 @@ public class AlignedFullPath implements IFullPath {
private final List<String> measurementList;
private final List<IMeasurementSchema> schemaList;
+ @Nullable private final Set<String> allSensors;
public AlignedFullPath(
IDeviceID deviceID, List<String> measurementList,
List<IMeasurementSchema> schemaList) {
this.deviceID = deviceID;
this.measurementList = measurementList;
this.schemaList = schemaList;
+ this.allSensors = null;
+ }
+
+ public AlignedFullPath(
+ IDeviceID deviceID,
+ List<String> measurementList,
+ List<IMeasurementSchema> schemaList,
+ Set<String> allSensors) {
+ this.deviceID = deviceID;
+ this.measurementList = measurementList;
+ this.schemaList = schemaList;
+ this.allSensors = allSensors;
}
@Override
@@ -68,6 +85,17 @@ public class AlignedFullPath implements IFullPath {
return measurementList.size();
}
+ public Set<String> getAllSensors() {
+ if (allSensors != null) {
+ return allSensors;
+ } else {
+ Set<String> res = new HashSet<>(measurementList);
+ // for time column
+ res.add("");
+ return res;
+ }
+ }
+
@Override
public long ramBytesUsed() {
return INSTANCE_SIZE