This is an automated email from the ASF dual-hosted git repository. caogaofei pushed a commit to branch fix_subquery in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 115ca1c50906004efcec5f53861adb2980d14165 Author: Beyyes <[email protected]> AuthorDate: Thu Oct 24 19:09:50 2024 +0800 fix subquery --- .../db/it/IoTDBMultiIDsWithAttributesTableIT.java | 39 +++++++++++++++++----- .../TableAggregationTableScanOperator.java | 36 +++++++++++++------- .../TableModelStatementMemorySourceVisitor.java | 6 ++-- .../plan/planner/OperatorTreeGenerator.java | 2 +- .../plan/planner/plan/node/PlanGraphPrinter.java | 2 +- .../plan/relational/analyzer/Analysis.java | 12 +++++++ .../distribute/TableDistributedPlanGenerator.java | 24 ++++++++----- .../optimizations/PushPredicateIntoTableScan.java | 8 +++-- 8 files changed, 94 insertions(+), 35 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBMultiIDsWithAttributesTableIT.java b/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBMultiIDsWithAttributesTableIT.java index ec659eddb42..5f1c6e2437c 100644 --- a/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBMultiIDsWithAttributesTableIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/relational/it/db/it/IoTDBMultiIDsWithAttributesTableIT.java @@ -899,14 +899,15 @@ public class IoTDBMultiIDsWithAttributesTableIT { }; // TODO(beyyes) test below - // sql = "select count(*) from (\n" + - // "\tselect device, level, date_bin(1d, time) as bin, \n" + - // "\tcount(num) as count_num, count(*) as count_star, count(device) as count_device, - // count(date) as count_date, count(attr1) as count_attr1, count(attr2) as count_attr2, - // count(time) as count_time, avg(num) as avg_num \n" + - // "\tfrom table0 \n" + - // "\tgroup by 3, device, level order by device, level, bin\n" + - // ")\n"; + // sql = "select count(*) from (\n" + + // "\tselect device, level, date_bin(1d, time) as bin, \n" + + // "\tcount(num) as count_num, count(*) as count_star, count(device) as + // count_device, + // count(date) as count_date, count(attr1) as count_attr1, count(attr2) as count_attr2, + // count(time) as count_time, avg(num) as avg_num \n" + + // "\tfrom table0 \n" + + // "\tgroup by 3, device, level order by device, level, bin\n" + + // ")\n"; } @Test @@ -1004,6 +1005,23 @@ public class IoTDBMultiIDsWithAttributesTableIT { + "count(attr1) as count_attr1, count(attr2) as count_attr2, count(time) as count_time, sum(num) as sum_num," + "avg(num) as avg_num from table0 where time=32 or time=1971-04-27T01:46:40.000+08:00 group by 3, device, level order by device, level"; tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME); + + // not exist device test + expectedHeader = buildHeaders(3); + sql = "select count(*), count(num), sum(num) from table0 where device='d_not_exist'"; + retArray = new String[] {"0,0,null,"}; + tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME); + + // not exist time range test + sql = "select count(*), count(num), sum(num) from table0 where time>2100-04-26T18:01:40.000"; + retArray = new String[] {"0,0,null,"}; + tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME); + + // only one device in this time + expectedHeader = buildHeaders(2); + sql = "select count(num),sum(num) from table1 where time=0"; + retArray = new String[] {"2,6.0,"}; + tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME); } @Test @@ -1150,6 +1168,11 @@ public class IoTDBMultiIDsWithAttributesTableIT { "1970-01-01T00:00:00.100Z,1970-01-01T00:00:00.100Z,1970-01-01T00:00:00.100Z,1971-01-01T00:00:00.000Z,1971-01-01T00:00:00.000Z,1970-01-01T00:00:00.100Z,1970-01-01T00:00:00.100Z,1970-01-01T00:00:00.100Z,1970-01-01T00:00:00.100Z,1970-01-01T00:00:00.100Z,1971-08-20T11:33:20.000Z,1971-01-01T00:01:40.000Z,1971-01-01T00:01:40.000Z,", }; tableResultSetEqualTest(sql, expectedHeader1, retArray, DATABASE_NAME); + + expectedHeader = buildHeaders(3); + sql = "select last_by(time,num+1),last_by(num+1,time),last_by(num+1,floatnum+1) from table0"; + retArray = new String[] {"1971-08-20T11:33:20.000Z,16,16,"}; + tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME); } @Test 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 a96240f62b7..51a42e4ce3a 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 @@ -40,6 +40,7 @@ import org.apache.iotdb.db.storageengine.dataregion.read.QueryDataSource; import org.apache.tsfile.block.column.Column; import org.apache.tsfile.block.column.ColumnBuilder; import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.file.metadata.StringArrayDeviceID; import org.apache.tsfile.file.metadata.statistics.Statistics; import org.apache.tsfile.file.metadata.statistics.StringStatistics; import org.apache.tsfile.read.common.TimeRange; @@ -55,6 +56,7 @@ import org.apache.tsfile.write.schema.IMeasurementSchema; import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Optional; import java.util.concurrent.TimeUnit; @@ -170,7 +172,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation this.maxReturnSize = maxReturnSize; this.maxTsBlockLineNum = maxTsBlockLineNum; - this.seriesScanUtil = constructAlignedSeriesScanUtil(deviceEntries.get(currentDeviceIndex)); + constructAlignedSeriesScanUtil(); } @Override @@ -253,17 +255,28 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation return resultTsBlock; } - private AlignedSeriesScanUtil constructAlignedSeriesScanUtil(DeviceEntry deviceEntry) { + private void constructAlignedSeriesScanUtil() { + DeviceEntry deviceEntry; + + if (this.deviceEntries.size() <= this.currentDeviceIndex + || this.deviceEntries.get(this.currentDeviceIndex) == null) { + // for device which not exist + deviceEntry = new DeviceEntry(new StringArrayDeviceID(""), Collections.emptyList()); + } else { + deviceEntry = this.deviceEntries.get(this.currentDeviceIndex); + } + AlignedFullPath alignedPath = constructAlignedPath(deviceEntry, measurementColumnNames, measurementSchemas); - return new AlignedSeriesScanUtil( - alignedPath, - scanOrder, - seriesScanOptions, - operatorContext.getInstanceContext(), - true, - measurementColumnTSDataTypes); + this.seriesScanUtil = + new AlignedSeriesScanUtil( + alignedPath, + scanOrder, + seriesScanOptions, + operatorContext.getInstanceContext(), + true, + measurementColumnTSDataTypes); } /** Return true if we have the result of this timeRange. */ @@ -313,7 +326,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation if (currentDeviceIndex < deviceCount) { // construct AlignedSeriesScanUtil for next device - this.seriesScanUtil = constructAlignedSeriesScanUtil(deviceEntries.get(currentDeviceIndex)); + constructAlignedSeriesScanUtil(); queryDataSource.reset(); this.seriesScanUtil.initQueryDataSource(queryDataSource); } @@ -790,8 +803,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation if (currentDeviceIndex < deviceCount) { // construct AlignedSeriesScanUtil for next device - this.seriesScanUtil = - constructAlignedSeriesScanUtil(deviceEntries.get(currentDeviceIndex)); + constructAlignedSeriesScanUtil(); queryDataSource.reset(); this.seriesScanUtil.initQueryDataSource(queryDataSource); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/memory/TableModelStatementMemorySourceVisitor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/memory/TableModelStatementMemorySourceVisitor.java index ec603a90e83..e87f2d24a2b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/memory/TableModelStatementMemorySourceVisitor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/memory/TableModelStatementMemorySourceVisitor.java @@ -75,9 +75,9 @@ public class TableModelStatementMemorySourceVisitor symbolAllocator, NOOP) .plan(context.getAnalysis()); - if (context.getAnalysis().isEmptyDataSource()) { - return new StatementMemorySource(new TsBlock(0), header); - } + // if (context.getAnalysis().isEmptyDataSource()) { + // return new StatementMemorySource(new TsBlock(0), header); + // } // Generate table model distributed plan final TableDistributedPlanGenerator.PlanContext planContext = 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 5dd0c61bd75..3f70b46398c 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 @@ -671,7 +671,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP boolean canUseStatistics = !TSDataType.BLOB.equals(node.getSeriesPath().getSeriesType()) || (aggregationDescriptors.stream() - .noneMatch(o -> !judgeCanUseStatistics(o.getAggregationType(), TSDataType.BLOB))); + .allMatch(o -> judgeCanUseStatistics(o.getAggregationType(), TSDataType.BLOB))); SeriesAggregationScanOperator aggregateScanOperator = new SeriesAggregationScanOperator( node.getPlanNodeId(), 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 84b28912569..d33f0b24481 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 @@ -696,7 +696,7 @@ public class PlanGraphPrinter extends PlanVisitor<List<String>, PlanGraphPrinter String.format( "RegionId: %s", node.getRegionReplicaSet() == null || node.getRegionReplicaSet().getRegionId() == null - ? "" + ? "Not Assigned" : node.getRegionReplicaSet().getRegionId().getId())); return render(node, boxValue, context); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/analyzer/Analysis.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/analyzer/Analysis.java index 32f0e5aec92..1b34213f96c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/analyzer/Analysis.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/analyzer/Analysis.java @@ -20,6 +20,7 @@ package org.apache.iotdb.db.queryengine.plan.relational.analyzer; import org.apache.iotdb.common.rpc.thrift.TEndPoint; +import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.commons.partition.DataPartition; import org.apache.iotdb.commons.partition.SchemaPartition; @@ -68,8 +69,10 @@ import com.google.common.collect.ImmutableSet; import com.google.common.collect.Multimap; import com.google.common.collect.Streams; import com.google.errorprone.annotations.Immutable; +import org.apache.tsfile.file.metadata.IDeviceID; import org.apache.tsfile.read.common.block.TsBlock; import org.apache.tsfile.read.common.type.Type; +import org.apache.tsfile.read.filter.basic.Filter; import org.apache.tsfile.utils.TimeDuration; import javax.annotation.Nullable; @@ -809,6 +812,15 @@ public class Analysis implements IAnalysis { redirectNodeList.add(endPoint); } + public List<TRegionReplicaSet> getDataRegionReplicaSetWithTimeFilter( + String database, IDeviceID deviceId, Filter timeFilter) { + if (dataPartition == null) { + return Collections.singletonList(new TRegionReplicaSet()); + } else { + return dataPartition.getDataRegionReplicaSetWithTimeFilter(database, deviceId, timeFilter); + } + } + @Override public TimePredicate getCovertedTimePredicate() { return null; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java index 848dbf4c917..7b35709522e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java @@ -571,17 +571,21 @@ public class TableDistributedPlanGenerator List<List<TRegionReplicaSet>> regionReplicaSetsList = new ArrayList<>(); for (DeviceEntry deviceEntry : node.getDeviceEntries()) { List<TRegionReplicaSet> regionReplicaSets = - analysis - .getDataPartitionInfo() - .getDataRegionReplicaSetWithTimeFilter( - node.getQualifiedObjectName().getDatabaseName(), - deviceEntry.getDeviceID(), - node.getTimeFilter()); + analysis.getDataRegionReplicaSetWithTimeFilter( + node.getQualifiedObjectName().getDatabaseName(), + deviceEntry.getDeviceID(), + node.getTimeFilter()); if (regionReplicaSets.size() > 1) { needSplit = true; } regionReplicaSetsList.add(regionReplicaSets); } + + if (regionReplicaSetsList.isEmpty()) { + regionReplicaSetsList = + Collections.singletonList(Collections.singletonList(new TRegionReplicaSet())); + } + // Step is SINGLE, has date_bin(time) and device data in more than one region, we need to split // this node into two-stage Aggregation needSplit = needSplit && node.getProjection() != null && node.getStep() == SINGLE; @@ -620,7 +624,9 @@ public class TableDistributedPlanGenerator scanNode.setRegionReplicaSet(regionReplicaSet); return scanNode; }); - aggregationTableScanNode.appendDeviceEntry(node.getDeviceEntries().get(i)); + if (node.getDeviceEntries().get(i) != null) { + aggregationTableScanNode.appendDeviceEntry(node.getDeviceEntries().get(i)); + } } } } else { @@ -653,7 +659,9 @@ public class TableDistributedPlanGenerator scanNode.setRegionReplicaSet(regionReplicaSet); return scanNode; }); - aggregationTableScanNode.appendDeviceEntry(node.getDeviceEntries().get(i)); + if (node.getDeviceEntries().size() > i && node.getDeviceEntries().get(i) != null) { + aggregationTableScanNode.appendDeviceEntry(node.getDeviceEntries().get(i)); + } } } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PushPredicateIntoTableScan.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PushPredicateIntoTableScan.java index 9db2f87cfd9..e3d9e6b3f84 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PushPredicateIntoTableScan.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/PushPredicateIntoTableScan.java @@ -483,7 +483,9 @@ public class PushPredicateIntoTableScan implements PlanOptimizer { .recordPlanCost(TABLE_TYPE, SCHEMA_FETCHER, System.nanoTime() - startTime); if (deviceEntries.isEmpty()) { - analysis.setFinishQueryAfterAnalyze(); + if (!analysis.hasAggregates()) { + analysis.setFinishQueryAfterAnalyze(); + } analysis.setEmptyDataSource(true); } else { Filter timeFilter = @@ -506,7 +508,9 @@ public class PushPredicateIntoTableScan implements PlanOptimizer { } if (dataPartition.getDataPartitionMap().isEmpty()) { - analysis.setFinishQueryAfterAnalyze(); + if (!analysis.hasAggregates()) { + analysis.setFinishQueryAfterAnalyze(); + } analysis.setEmptyDataSource(true); } else { analysis.upsertDataPartition(dataPartition);
