This is an automated email from the ASF dual-hosted git repository. caogaofei pushed a commit to branch fix_having_again in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 12e0473157c3506c6f39574dcc18519f94987969 Author: Beyyes <[email protected]> AuthorDate: Mon Oct 28 14:17:16 2024 +0800 temp --- .../db/it/IoTDBMultiIDsWithAttributesTableIT.java | 8 +- .../TableAggregationTableScanOperator.java | 117 +++++++++++---------- .../aggregation/TableModeAccumulator.java | 4 +- .../queryengine/plan/analyze/ExpressionUtils.java | 19 +++- .../plan/planner/TableOperatorGenerator.java | 28 ----- 5 files changed, 81 insertions(+), 95 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 1b4f100bb44..5d43ee7a4b3 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 @@ -47,7 +47,7 @@ public class IoTDBMultiIDsWithAttributesTableIT { "CREATE DATABASE db", "USE db", "CREATE TABLE table0 (device string id, level string id, attr1 string attribute, attr2 string attribute, num int32 measurement, bigNum int64 measurement, " - + "floatNum FLOAT measurement, str TEXT measurement, bool BOOLEAN measurement, date DATE measurement, blob BLOB measurement, ts TIMESTAMP measurement, stringV STRING measurement, doubleNum DOUBLE measurement)", + + "floatNum float measurement, str TEXT measurement, bool BOOLEAN measurement, date DATE measurement, blob BLOB measurement, ts TIMESTAMP measurement, stringV STRING measurement, doubleNum DOUBLE measurement)", "insert into table0(device, level, attr1, attr2, time,num,bigNum,floatNum,str,bool) values('d1', 'l1', 'c', 'd', 0,3,2947483648,231.2121,'coconut',FALSE)", "insert into table0(device, level, attr1, attr2, time,num,bigNum,floatNum,str,bool,blob,ts,doubleNum) values('d1', 'l2', 'y', 'z', 20,2,2147483648,434.12,'pineapple',TRUE,X'108DCD62',2024-09-24T06:15:35.000+00:00,6666.8)", "insert into table0(device, level, attr1, attr2, time,num,bigNum,floatNum,str,bool) values('d1', 'l3', 't', 'a', 40,1,2247483648,12.123,'apricot',TRUE)", @@ -901,12 +901,12 @@ public class IoTDBMultiIDsWithAttributesTableIT { tableResultSetEqualTest(sql, expectedHeader, retArray, DATABASE_NAME); // flush multi times, generated multi tsfile - expectedHeader = buildHeaders(1); - sql = "select date_bin(40ms,time), first(time) from table1 where device='d11' group by 1"; + expectedHeader = buildHeaders(2); + sql = "select date_bin(30ms,time), first(time) from table1 where device='d11' group by 1"; retArray = new String[] { "1970-01-01T00:00:00.000Z,1970-01-01T00:00:00.000Z,", - "1970-01-01T00:00:00.000Z,1970-01-01T00:00:00.000Z," + "1970-01-01T00:00:00.030Z,1970-01-01T00:00:00.030Z," }; // TODO(beyyes) test below 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..fe086fd9e13 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 @@ -24,7 +24,7 @@ import org.apache.iotdb.commons.schema.table.column.TsTableColumnCategory; import org.apache.iotdb.db.queryengine.execution.MemoryEstimationHelper; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITableTimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; -import org.apache.iotdb.db.queryengine.execution.operator.source.AbstractSeriesAggregationScanOperator; +import org.apache.iotdb.db.queryengine.execution.operator.source.AbstractDataSourceOperator; import org.apache.iotdb.db.queryengine.execution.operator.source.AlignedSeriesScanUtil; import org.apache.iotdb.db.queryengine.execution.operator.source.relational.aggregation.TableAggregator; import org.apache.iotdb.db.queryengine.execution.operator.window.IWindow; @@ -39,6 +39,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.common.conf.TSFileDescriptor; import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.file.metadata.StringArrayDeviceID; import org.apache.tsfile.file.metadata.statistics.Statistics; @@ -67,53 +68,54 @@ import static org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil import static org.apache.iotdb.db.queryengine.execution.operator.source.relational.TableScanOperator.constructAlignedPath; import static org.apache.tsfile.read.common.block.TsBlockUtil.skipPointsOutOfTimeRange; -public class TableAggregationTableScanOperator extends AbstractSeriesAggregationScanOperator { +public class TableAggregationTableScanOperator extends AbstractDataSourceOperator { + + public static final LongColumn TIME_COLUMN_TEMPLATE = + new LongColumn(1, Optional.empty(), new long[] {0}); private static final long INSTANCE_SIZE = RamUsageEstimator.shallowSizeOfInstance(TableAggregationTableScanOperator.class); + // can not calc maxTsBlockLineNum using date_bin + private final int maxTsBlockLineNum = + TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(); + private final long cachedRawDataSize; - private final List<TableAggregator> tableAggregators; + private boolean finished = false; + private final List<TableAggregator> tableAggregators; private final List<ColumnSchema> groupingKeySchemas; private final int[] groupingKeyIndex; - - public static final LongColumn TIME_COLUMN_TEMPLATE = - new LongColumn(1, Optional.empty(), new long[] {0}); + // for different aggregations aiming to same column, use this variable to point to same column + private final List<Integer> aggArguments; private final List<ColumnSchema> columnSchemas; - private final int[] columnsIndexArray; private final List<DeviceEntry> deviceEntries; - private final int deviceCount; - - private final Ordering scanOrder; - private final SeriesScanOptions seriesScanOptions; + private final int measurementCount; + private int currentDeviceIndex; private final List<String> measurementColumnNames; - private final List<IMeasurementSchema> measurementSchemas; - private final List<TSDataType> measurementColumnTSDataTypes; - // TODO calc maxTsBlockLineNum using date_bin - private final int maxTsBlockLineNum; - - // for different aggregations aiming to same column, use this variable to point to same column - private final List<Integer> aggArguments; - + private final Ordering scanOrder; + private final SeriesScanOptions seriesScanOptions; private QueryDataSource queryDataSource; - private int currentDeviceIndex; - - ITableTimeRangeIterator timeIterator; + private final ITableTimeRangeIterator timeIterator; + private final boolean canUseStatistics; + private final boolean ascending; private boolean allAggregatorsHasFinalResult = false; + private long leftRuntimeOfOneNextCall; + + private TsBlock inputTsBlock; public TableAggregationTableScanOperator( PlanNodeId sourceId, - OperatorContext context, + OperatorContext operatorContext, List<ColumnSchema> columnSchemas, int[] columnsIndexArray, List<DeviceEntry> deviceEntries, @@ -121,36 +123,25 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation SeriesScanOptions seriesScanOptions, List<String> measurementColumnNames, List<IMeasurementSchema> measurementSchemas, - int maxTsBlockLineNum, int measurementCount, List<TableAggregator> tableAggregators, List<ColumnSchema> groupingKeySchemas, int[] groupingKeyIndex, ITableTimeRangeIterator tableTimeRangeIterator, boolean ascending, - long maxReturnSize, boolean canUseStatistics, List<Integer> aggArguments) { - super( - sourceId, - context, - null, - measurementCount, - null, - null, - ascending, - false, - null, - maxReturnSize, - canUseStatistics); + this.sourceId = sourceId; + this.operatorContext = operatorContext; + this.canUseStatistics = canUseStatistics; + this.ascending = ascending; + this.measurementCount = measurementCount; this.tableAggregators = tableAggregators; this.groupingKeySchemas = groupingKeySchemas; this.groupingKeyIndex = groupingKeyIndex; - this.sourceId = sourceId; - this.operatorContext = context; this.columnSchemas = columnSchemas; this.columnsIndexArray = columnsIndexArray; this.deviceEntries = deviceEntries; @@ -164,13 +155,10 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation this.currentDeviceIndex = 0; this.aggArguments = aggArguments; this.timeIterator = tableTimeRangeIterator; - if (tableTimeRangeIterator.getType() - == ITableTimeRangeIterator.TimeIteratorType.SINGLE_TIME_ITERATOR) { - curTimeRange = new TimeRange(Long.MIN_VALUE, Long.MAX_VALUE); - } - this.maxReturnSize = maxReturnSize; - this.maxTsBlockLineNum = maxTsBlockLineNum; + this.maxReturnSize = TSFileDescriptor.getInstance().getConfig().getMaxTsBlockSizeInBytes(); + this.cachedRawDataSize = + (1L + measurementCount) * TSFileDescriptor.getInstance().getConfig().getPageSizeInByte(); constructAlignedSeriesScanUtil(); } @@ -279,7 +267,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation } /** Return true if we have the result of this timeRange. */ - @Override + // @Override protected boolean calculateAggregationResultForCurrentTimeRange() { try { if (calcFromCachedData()) { @@ -342,14 +330,14 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation } } - @Override + // @Override protected void updateResultTsBlock() { appendAggregationResult(resultTsBlockBuilder, tableAggregators); // after appendAggregationResult invoked, aggregators must be cleared resetTableAggregators(); } - @Override + // @Override protected boolean calcFromCachedData() { return calcUsingRawData(inputTsBlock); } @@ -552,7 +540,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation } @SuppressWarnings({"squid:S3776", "squid:S135", "squid:S3740"}) - @Override + // @Override public boolean readAndCalcFromFile() throws IOException { // start stopwatch long start = System.nanoTime(); @@ -575,8 +563,8 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation if (timeIterator .getCurTimeRange() .contains(fileTimeStatistics.getStartTime(), fileTimeStatistics.getEndTime())) { - Statistics[] statisticsList = new Statistics[subSensorSize]; - for (int i = 0; i < subSensorSize; i++) { + Statistics[] statisticsList = new Statistics[measurementCount]; + for (int i = 0; i < measurementCount; i++) { statisticsList[i] = seriesScanUtil.currentFileStatistics(i); } calcFromStatistics(fileTimeStatistics, statisticsList); @@ -622,8 +610,8 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation .getCurTimeRange() .contains(chunkTimeStatistics.getStartTime(), chunkTimeStatistics.getEndTime())) { // calc from chunkMetaData - Statistics[] statisticsList = new Statistics[subSensorSize]; - for (int i = 0; i < subSensorSize; i++) { + Statistics[] statisticsList = new Statistics[measurementCount]; + for (int i = 0; i < measurementCount; i++) { statisticsList[i] = seriesScanUtil.currentChunkStatistics(i); } calcFromStatistics(chunkTimeStatistics, statisticsList); @@ -644,8 +632,6 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation return false; } - long leftRuntimeOfOneNextCall = Long.MAX_VALUE; - @SuppressWarnings({"squid:S3776", "squid:S135", "squid:S3740"}) protected boolean readAndCalcFromPage() throws IOException { long start = System.nanoTime(); @@ -671,8 +657,8 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation if (timeIterator .getCurTimeRange() .contains(pageTimeStatistics.getStartTime(), pageTimeStatistics.getEndTime())) { - Statistics[] statisticsList = new Statistics[subSensorSize]; - for (int i = 0; i < subSensorSize; i++) { + Statistics[] statisticsList = new Statistics[measurementCount]; + for (int i = 0; i < measurementCount; i++) { statisticsList[i] = seriesScanUtil.currentPageStatistics(i); } calcFromStatistics(pageTimeStatistics, statisticsList); @@ -863,4 +849,21 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation + (resultTsBlockBuilder == null ? 0 : resultTsBlockBuilder.getRetainedSizeInBytes()) + RamUsageEstimator.sizeOfCollection(deviceEntries); } + + @Override + public long calculateMaxPeekMemory() { + return 0; + } + + @Override + public long calculateMaxReturnSize() { + return 0; + } + + @Override + public long calculateRetainedSizeAfterCallingNext() { + return timeIterator.getType() == ITableTimeRangeIterator.TimeIteratorType.DATE_BIN_TIME_ITERATOR + ? cachedRawDataSize + : 0; + } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/TableModeAccumulator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/TableModeAccumulator.java index 1ad01095fac..66ddece337b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/TableModeAccumulator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/TableModeAccumulator.java @@ -531,11 +531,11 @@ public class TableModeAccumulator implements TableAccumulator { } private void checkMapSize(int size) { - if (size > MAP_SIZE_THRESHOLD) { + if (size > IoTDBDescriptor.getInstance().getConfig().getModeMapSizeThreshold()) { throw new RuntimeException( String.format( "distinct values has exceeded the threshold %s when calculate Mode", - MAP_SIZE_THRESHOLD)); + IoTDBDescriptor.getInstance().getConfig().getModeMapSizeThreshold())); } } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ExpressionUtils.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ExpressionUtils.java index 6d76384c1c8..e8eb2d69760 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ExpressionUtils.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/ExpressionUtils.java @@ -200,11 +200,22 @@ public class ExpressionUtils { final List<Expression> rightExpressions, final MPPQueryContext queryContext) { List<Expression> resultExpressions = new ArrayList<>(); - for (Expression le : leftExpressions) { + if (!leftExpressions.isEmpty() && !rightExpressions.isEmpty()) { + for (Expression le : leftExpressions) { + for (Expression re : rightExpressions) { + resultExpressions.add( + reserveMemoryForExpression( + queryContext, reconstructBinaryExpression(expression, le, re))); + } + } + return resultExpressions; + } else if (!leftExpressions.isEmpty()) { + for (Expression le : leftExpressions) { + resultExpressions.add(reserveMemoryForExpression(queryContext, le)); + } + } else if (!rightExpressions.isEmpty()) { for (Expression re : rightExpressions) { - resultExpressions.add( - reserveMemoryForExpression( - queryContext, reconstructBinaryExpression(expression, le, re))); + resultExpressions.add(reserveMemoryForExpression(queryContext, re)); } } return resultExpressions; 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 56d6868a24b..bcf320fe76d 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 @@ -1643,14 +1643,12 @@ public class TableOperatorGenerator extends PlanVisitor<Operator, LocalExecution scanOptionsBuilder.build(), measurementColumnNames, measurementSchemas, - TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(), measurementColumnCount, aggregators, groupingKeySchemas, groupingKeyIndex, timeRangeIterator, scanAscending, - calculateMaxAggregationResultSize(), canUseStatistic, aggColumnIndexes); @@ -1734,30 +1732,4 @@ public class TableOperatorGenerator extends PlanVisitor<Operator, LocalExecution } return new boolean[] {canUseStatistic, isAscending}; } - - public static long calculateMaxAggregationResultSize( - // List<? extends AggregationDescriptor> aggregationDescriptors, - // ITimeRangeIterator timeRangeIterator - ) { - // TODO perfect max aggregation result size logic - return TSFileDescriptor.getInstance().getConfig().getMaxTsBlockSizeInBytes(); - - // long timeValueColumnsSizePerLine = TimeColumn.SIZE_IN_BYTES_PER_POSITION; - // for (AggregationDescriptor descriptor : aggregationDescriptors) { - // List<TSDataType> outPutDataTypes = - // descriptor.getOutputColumnNames().stream() - // .map(typeProvider::getTableModelType) - // .collect(Collectors.toList()); - // for (TSDataType tsDataType : outPutDataTypes) { - // timeValueColumnsSizePerLine += getOutputColumnSizePerLine(tsDataType); - // } - // } - // - // return Math.min( - // TSFileDescriptor.getInstance().getConfig().getMaxTsBlockSizeInBytes(), - // Math.min( - // TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(), - // timeRangeIterator.getTotalIntervalNum()) - // * timeValueColumnsSizePerLine); - } }
