This is an automated email from the ASF dual-hosted git repository. caogaofei pushed a commit to branch agg_table_scan in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 22f30d046f151c21226c32bdfdd58a132a30f96a Author: Beyyes <[email protected]> AuthorDate: Sun Sep 29 20:07:22 2024 +0800 tmp --- .../protocol/thrift/impl/ClientRPCServiceImpl.java | 6 +- .../execution/aggregation/IAggregator.java | 45 ++++ .../{Aggregator.java => TreeAggregator.java} | 6 +- .../slidingwindow/SlidingWindowAggregator.java | 4 +- .../execution/operator/AggregationUtil.java | 23 +- .../operator/process/AggregationOperator.java | 12 +- .../process/RawDataAggregationOperator.java | 8 +- .../process/SingleInputAggregationOperator.java | 6 +- .../process/SlidingWindowAggregationOperator.java | 10 +- .../operator/process/TagAggregationOperator.java | 21 +- .../operator/process/last/LastQueryUtil.java | 18 +- .../AbstractSeriesAggregationScanOperator.java | 12 +- .../AlignedSeriesAggregationScanOperator.java | 6 +- .../source/SeriesAggregationScanOperator.java | 6 +- .../TableAggregationTableScanOperator.java | 289 ++++++++++++++++++++- .../aggregation/AggregationOperator.java | 10 +- .../{Aggregator.java => TableAggregator.java} | 4 +- .../operator/window/ConditionWindowManager.java | 6 +- .../operator/window/CountWindowManager.java | 6 +- .../execution/operator/window/IWindowManager.java | 15 +- .../operator/window/SessionWindowManager.java | 6 +- .../operator/window/TimeWindowManager.java | 6 +- .../operator/window/VariationWindowManager.java | 6 +- .../plan/planner/OperatorTreeGenerator.java | 36 +-- .../plan/planner/TableOperatorGenerator.java | 27 +- .../operator/AggregationOperatorTest.java | 10 +- .../AlignedSeriesAggregationScanOperatorTest.java | 73 +++--- .../operator/HorizontallyConcatOperatorTest.java | 6 +- .../execution/operator/LastQueryOperatorTest.java | 10 +- .../operator/LastQueryTreeSortOperatorTest.java | 10 +- .../execution/operator/OperatorMemoryTest.java | 18 +- .../operator/RawDataAggregationOperatorTest.java | 6 +- .../SeriesAggregationScanOperatorTest.java | 68 ++--- .../SlidingWindowAggregationOperatorTest.java | 9 +- .../operator/UpdateLastCacheOperatorTest.java | 10 +- 35 files changed, 574 insertions(+), 240 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java index 4462ad66fce..18363ce11e7 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/ClientRPCServiceImpl.java @@ -63,7 +63,7 @@ import org.apache.iotdb.db.queryengine.common.header.ColumnHeader; import org.apache.iotdb.db.queryengine.common.header.DatasetHeader; import org.apache.iotdb.db.queryengine.common.header.DatasetHeaderFactory; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceManager; @@ -752,8 +752,8 @@ public class ClientRPCServiceImpl implements IClientRPCServiceWithHandler { scanOptionsBuilder.withGlobalTimeFilter(timeFilter); String aggregationName = SchemaUtils.getBuiltinAggregationName(aggregationType); - Aggregator aggregator = - new Aggregator( + TreeAggregator aggregator = + new TreeAggregator( AccumulatorFactory.createAccumulator( aggregationName, aggregationType, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/IAggregator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/IAggregator.java new file mode 100644 index 00000000000..1cf694d2ba4 --- /dev/null +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/IAggregator.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.queryengine.execution.aggregation; + +import org.apache.tsfile.block.column.ColumnBuilder; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.file.metadata.statistics.Statistics; +import org.apache.tsfile.read.common.block.TsBlock; +import org.apache.tsfile.utils.BitMap; + +public interface IAggregator { + + // Tree model: used for SeriesAggregateScanOperator and RawDataAggregateOperator + void processTsBlock(TsBlock tsBlock, BitMap bitMap); + + // Tree model: used for AggregateOperator + void processTsBlocks(TsBlock[] tsBlock); + + void outputResult(ColumnBuilder[] columnBuilder); + + void processStatistics(Statistics timeStatistics, Statistics[] valueStatistics); + + TSDataType[] getOutputType(); + + void reset(); + + boolean hasFinalResult(); +} diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/Aggregator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/TreeAggregator.java similarity index 97% rename from iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/Aggregator.java rename to iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/TreeAggregator.java index 6992d80b53f..769a878b32e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/Aggregator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/TreeAggregator.java @@ -37,7 +37,7 @@ import static com.google.common.base.Preconditions.checkArgument; import static org.apache.iotdb.db.queryengine.metric.QueryExecutionMetricSet.AGGREGATION_FROM_RAW_DATA; import static org.apache.iotdb.db.queryengine.metric.QueryExecutionMetricSet.AGGREGATION_FROM_STATISTICS; -public class Aggregator { +public class TreeAggregator implements IAggregator { protected final Accumulator accumulator; // In some intermediate result input, inputLocation[] should include two columns @@ -47,7 +47,7 @@ public class Aggregator { QueryExecutionMetricSet.getInstance(); // Used for SeriesAggregateScanOperator - public Aggregator(Accumulator accumulator, AggregationStep step) { + public TreeAggregator(Accumulator accumulator, AggregationStep step) { this.accumulator = accumulator; this.step = step; this.inputLocationList = @@ -55,7 +55,7 @@ public class Aggregator { } // Used for AggregateOperator, AlignedSeriesAggregateScanOperator - public Aggregator( + public TreeAggregator( Accumulator accumulator, AggregationStep step, List<InputLocation[]> inputLocationList) { this.accumulator = accumulator; this.step = step; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/slidingwindow/SlidingWindowAggregator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/slidingwindow/SlidingWindowAggregator.java index 3421487700e..adc179f77ec 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/slidingwindow/SlidingWindowAggregator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/aggregation/slidingwindow/SlidingWindowAggregator.java @@ -20,7 +20,7 @@ package org.apache.iotdb.db.queryengine.execution.aggregation.slidingwindow; import org.apache.iotdb.db.queryengine.execution.aggregation.Accumulator; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.plan.planner.plan.parameter.AggregationStep; import org.apache.iotdb.db.queryengine.plan.planner.plan.parameter.InputLocation; @@ -37,7 +37,7 @@ import java.util.stream.Collectors; import static com.google.common.base.Preconditions.checkArgument; -public abstract class SlidingWindowAggregator extends Aggregator { +public abstract class SlidingWindowAggregator extends TreeAggregator { // cached partial aggregation result of pre-aggregate windows protected Deque<PartialAggregationResult> deque; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationUtil.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationUtil.java index cbf30889650..6bcfa9a3169 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationUtil.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationUtil.java @@ -20,7 +20,7 @@ package org.apache.iotdb.db.queryengine.execution.operator; import org.apache.iotdb.commons.udf.builtin.BuiltinAggregationFunction; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.SingleTimeWindowIterator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.TimeRangeIteratorFactory; @@ -96,7 +96,7 @@ public class AggregationUtil { */ public static Pair<Boolean, TsBlock> calculateAggregationFromRawData( TsBlock inputTsBlock, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, TimeRange curTimeRange, boolean ascending) { if (inputTsBlock == null || inputTsBlock.isEmpty()) { @@ -124,8 +124,8 @@ public class AggregationUtil { isAllAggregatorsHasFinalResult(aggregators) || isTsBlockOutOfBound, inputTsBlock); } - private static TsBlock process( - TsBlock inputTsBlock, TimeRange curTimeRange, List<Aggregator> aggregators) { + public static TsBlock process( + TsBlock inputTsBlock, TimeRange curTimeRange, List<TreeAggregator> aggregators) { // Get the row which need to be processed by aggregator IWindow curWindow = new TimeWindow(curTimeRange); Column timeColumn = inputTsBlock.getTimeColumn(); @@ -138,7 +138,7 @@ public class AggregationUtil { } TsBlock inputRegion = inputTsBlock.getRegion(0, lastIndexToProcess + 1); - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { // current agg method has been calculated if (aggregator.hasFinalResult()) { continue; @@ -156,7 +156,7 @@ public class AggregationUtil { /** Append a row of aggregation results to the result tsBlock. */ public static void appendAggregationResult( TsBlockBuilder tsBlockBuilder, - List<? extends Aggregator> aggregators, + List<? extends TreeAggregator> aggregators, long outputTime, long endTime) { TimeColumnBuilder timeColumnBuilder = tsBlockBuilder.getTimeColumnBuilder(); @@ -168,7 +168,7 @@ public class AggregationUtil { columnBuilders[columnIndex].writeLong(endTime); columnIndex++; } - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { ColumnBuilder[] columnBuilder = new ColumnBuilder[aggregator.getOutputType().length]; columnBuilder[0] = columnBuilders[columnIndex++]; if (columnBuilder.length > 1) { @@ -180,7 +180,7 @@ public class AggregationUtil { } public static void appendAggregationResult( - TsBlockBuilder tsBlockBuilder, List<? extends Aggregator> aggregators, long outputTime) { + TsBlockBuilder tsBlockBuilder, List<? extends TreeAggregator> aggregators, long outputTime) { appendAggregationResult(tsBlockBuilder, aggregators, outputTime, INVALID_END_TIME); } @@ -198,8 +198,8 @@ public class AggregationUtil { && tsBlock.getStartTime() >= curTimeRange.getMin()); } - public static boolean isAllAggregatorsHasFinalResult(List<Aggregator> aggregators) { - for (Aggregator aggregator : aggregators) { + public static boolean isAllAggregatorsHasFinalResult(List<TreeAggregator> aggregators) { + for (TreeAggregator aggregator : aggregators) { if (!aggregator.hasFinalResult()) { return false; } @@ -230,7 +230,8 @@ public class AggregationUtil { * timeValueColumnsSizePerLine); } - public static long calculateMaxAggregationResultSizeForLastQuery(List<Aggregator> aggregators) { + public static long calculateMaxAggregationResultSizeForLastQuery( + List<TreeAggregator> aggregators) { long timeValueColumnsSizePerLine = TimeColumn.SIZE_IN_BYTES_PER_POSITION; List<TSDataType> outPutDataTypes = aggregators.stream() diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggregationOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggregationOperator.java index ff1de243b27..1e933f4b415 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggregationOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/AggregationOperator.java @@ -20,7 +20,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.process; import org.apache.iotdb.db.queryengine.execution.MemoryEstimationHelper; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.Operator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; @@ -52,7 +52,7 @@ public class AggregationOperator extends AbstractConsumeAllOperator { // Current interval of aggregation window [curStartTime, curEndTime) private TimeRange curTimeRange; - private final List<Aggregator> aggregators; + private final List<TreeAggregator> aggregators; // Using for building result tsBlock private final TsBlockBuilder resultTsBlockBuilder; @@ -63,7 +63,7 @@ public class AggregationOperator extends AbstractConsumeAllOperator { public AggregationOperator( OperatorContext operatorContext, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, List<Operator> children, boolean outputEndTime, @@ -75,7 +75,7 @@ public class AggregationOperator extends AbstractConsumeAllOperator { if (outputEndTime) { dataTypes.add(TSDataType.INT64); } - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { dataTypes.addAll(Arrays.asList(aggregator.getOutputType())); } this.resultTsBlockBuilder = new TsBlockBuilder(dataTypes); @@ -128,7 +128,7 @@ public class AggregationOperator extends AbstractConsumeAllOperator { curTimeRange = timeRangeIterator.nextTimeRange(); // Clear previous aggregation result - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { aggregator.reset(); } } @@ -153,7 +153,7 @@ public class AggregationOperator extends AbstractConsumeAllOperator { private void calculateNextAggregationResult() { // Consume current input tsBlocks - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { aggregator.processTsBlocks(inputTsBlocks); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/RawDataAggregationOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/RawDataAggregationOperator.java index 69f67503cc4..e177834b8f5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/RawDataAggregationOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/RawDataAggregationOperator.java @@ -20,7 +20,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.process; import org.apache.iotdb.db.queryengine.execution.MemoryEstimationHelper; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.Operator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; @@ -69,7 +69,7 @@ public class RawDataAggregationOperator extends SingleInputAggregationOperator { public RawDataAggregationOperator( OperatorContext operatorContext, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, Operator child, boolean ascending, @@ -195,7 +195,7 @@ public class RawDataAggregationOperator extends SingleInputAggregationOperator { } TsBlock inputRegion = inputTsBlock.getRegion(0, lastIndexToProcess + 1); - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { // Current agg method has been calculated if (aggregator.hasFinalResult()) { continue; @@ -255,7 +255,7 @@ public class RawDataAggregationOperator extends SingleInputAggregationOperator { private void initWindowAndAggregators() { windowManager.initCurWindow(); - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { aggregator.reset(); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SingleInputAggregationOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SingleInputAggregationOperator.java index 50fdff2eee3..1f36586653e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SingleInputAggregationOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SingleInputAggregationOperator.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.process; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.operator.Operator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; @@ -39,7 +39,7 @@ public abstract class SingleInputAggregationOperator implements ProcessOperator protected TsBlock inputTsBlock; protected boolean canCallNext; - protected final List<Aggregator> aggregators; + protected final List<TreeAggregator> aggregators; // using for building result tsBlock protected TsBlockBuilder resultTsBlockBuilder; @@ -49,7 +49,7 @@ public abstract class SingleInputAggregationOperator implements ProcessOperator protected SingleInputAggregationOperator( OperatorContext operatorContext, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, Operator child, boolean ascending, long maxReturnSize) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SlidingWindowAggregationOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SlidingWindowAggregationOperator.java index 9a1de9522af..ee460aa5677 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SlidingWindowAggregationOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/SlidingWindowAggregationOperator.java @@ -20,7 +20,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.process; import org.apache.iotdb.db.queryengine.execution.MemoryEstimationHelper; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.slidingwindow.SlidingWindowAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.Operator; @@ -55,7 +55,7 @@ public class SlidingWindowAggregationOperator extends SingleInputAggregationOper public SlidingWindowAggregationOperator( OperatorContext operatorContext, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, Operator child, boolean ascending, @@ -71,7 +71,7 @@ public class SlidingWindowAggregationOperator extends SingleInputAggregationOper if (outputEndTime) { dataTypes.add(TSDataType.INT64); } - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { dataTypes.addAll(Arrays.asList(aggregator.getOutputType())); } this.resultTsBlockBuilder = new TsBlockBuilder(dataTypes); @@ -94,7 +94,7 @@ public class SlidingWindowAggregationOperator extends SingleInputAggregationOper curTimeRange = timeRangeIterator.nextTimeRange(); // Clear previous aggregation result - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { ((SlidingWindowAggregator) aggregator).updateTimeRange(curTimeRange); } } @@ -142,7 +142,7 @@ public class SlidingWindowAggregationOperator extends SingleInputAggregationOper return; } - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { ((SlidingWindowAggregator) aggregator).processTsBlock(inputTsBlock); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/TagAggregationOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/TagAggregationOperator.java index d7399fb4554..e1869013365 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/TagAggregationOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/TagAggregationOperator.java @@ -20,7 +20,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.process; import org.apache.iotdb.db.queryengine.execution.MemoryEstimationHelper; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.operator.Operator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; @@ -44,7 +44,7 @@ public class TagAggregationOperator extends AbstractConsumeAllOperator { private static final long INSTANCE_SIZE = RamUsageEstimator.shallowSizeOfInstance(TagAggregationOperator.class); private final List<List<String>> groups; - private final List<List<Aggregator>> groupedAggregators; + private final List<List<TreeAggregator>> groupedAggregators; // These fields record the to be consumed index of each tsBlock. private final int[] consumedIndices; @@ -55,7 +55,7 @@ public class TagAggregationOperator extends AbstractConsumeAllOperator { public TagAggregationOperator( OperatorContext operatorContext, List<List<String>> groups, - List<List<Aggregator>> groupedAggregators, + List<List<TreeAggregator>> groupedAggregators, List<Operator> children, long maxReturnSize) { super(operatorContext, children); @@ -68,8 +68,8 @@ public class TagAggregationOperator extends AbstractConsumeAllOperator { for (int outputColumnIdx = 0; outputColumnIdx < groupedAggregators.get(0).size(); outputColumnIdx++) { - for (List<Aggregator> aggregators : groupedAggregators) { - Aggregator aggregator = aggregators.get(outputColumnIdx); + for (List<TreeAggregator> aggregators : groupedAggregators) { + TreeAggregator aggregator = aggregators.get(outputColumnIdx); if (aggregator != null) { actualOutputColumnTypes.addAll(Arrays.asList(aggregator.getOutputType())); break; @@ -109,7 +109,7 @@ public class TagAggregationOperator extends AbstractConsumeAllOperator { rowBlocks[i] = inputTsBlocks[i].getRegion(consumedIndices[i], 1); } for (int groupIdx = 0; groupIdx < groups.size(); groupIdx++) { - List<Aggregator> aggregators = groupedAggregators.get(groupIdx); + List<TreeAggregator> aggregators = groupedAggregators.get(groupIdx); aggregate(aggregators, rowBlocks); List<String> group = groups.get(groupIdx); appendOneRow(rowBlocks, group, aggregators); @@ -121,8 +121,8 @@ public class TagAggregationOperator extends AbstractConsumeAllOperator { } } - private void aggregate(List<Aggregator> aggregators, TsBlock[] rowBlocks) { - for (Aggregator aggregator : aggregators) { + private void aggregate(List<TreeAggregator> aggregators, TsBlock[] rowBlocks) { + for (TreeAggregator aggregator : aggregators) { if (aggregator == null) { continue; } @@ -131,7 +131,8 @@ public class TagAggregationOperator extends AbstractConsumeAllOperator { } } - private void appendOneRow(TsBlock[] rowBlocks, List<String> group, List<Aggregator> aggregators) { + private void appendOneRow( + TsBlock[] rowBlocks, List<String> group, List<TreeAggregator> aggregators) { TimeColumnBuilder timeColumnBuilder = tsBlockBuilder.getTimeColumnBuilder(); timeColumnBuilder.writeLong(rowBlocks[0].getStartTime()); ColumnBuilder[] columnBuilders = tsBlockBuilder.getValueColumnBuilders(); @@ -144,7 +145,7 @@ public class TagAggregationOperator extends AbstractConsumeAllOperator { } } for (int i = 0; i < aggregators.size(); i++) { - Aggregator aggregator = aggregators.get(i); + TreeAggregator aggregator = aggregators.get(i); ColumnBuilder columnBuilder = columnBuilders[i + group.size()]; if (aggregator == null) { columnBuilder.appendNull(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/last/LastQueryUtil.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/last/LastQueryUtil.java index 98525791f6d..d68ae9b9e8b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/last/LastQueryUtil.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/process/last/LastQueryUtil.java @@ -20,9 +20,9 @@ package org.apache.iotdb.db.queryengine.execution.operator.process.last; import org.apache.iotdb.commons.conf.CommonDescriptor; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.LastValueDescAccumulator; import org.apache.iotdb.db.queryengine.execution.aggregation.MaxTimeDescAccumulator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.plan.planner.plan.parameter.AggregationStep; import org.apache.iotdb.db.queryengine.plan.planner.plan.parameter.InputLocation; @@ -130,33 +130,33 @@ public class LastQueryUtil { return filter == null || filter.satisfy(tvPair.getTimestamp(), tvPair.getValue().getValue()); } - public static List<Aggregator> createAggregators(TSDataType dataType) { + public static List<TreeAggregator> createAggregators(TSDataType dataType) { // max_time, last_value - List<Aggregator> aggregators = new ArrayList<>(2); + List<TreeAggregator> aggregators = new ArrayList<>(2); aggregators.add( - new Aggregator( + new TreeAggregator( new MaxTimeDescAccumulator(), AggregationStep.SINGLE, Collections.singletonList(new InputLocation[] {new InputLocation(0, 0)}))); aggregators.add( - new Aggregator( + new TreeAggregator( new LastValueDescAccumulator(dataType), AggregationStep.SINGLE, Collections.singletonList(new InputLocation[] {new InputLocation(0, 0)}))); return aggregators; } - public static List<Aggregator> createAggregators(TSDataType dataType, int valueColumnIndex) { + public static List<TreeAggregator> createAggregators(TSDataType dataType, int valueColumnIndex) { // max_time, last_value - List<Aggregator> aggregators = new ArrayList<>(2); + List<TreeAggregator> aggregators = new ArrayList<>(2); aggregators.add( - new Aggregator( + new TreeAggregator( new MaxTimeDescAccumulator(), AggregationStep.SINGLE, Collections.singletonList( new InputLocation[] {new InputLocation(0, valueColumnIndex)}))); aggregators.add( - new Aggregator( + new TreeAggregator( new LastValueDescAccumulator(dataType), AggregationStep.SINGLE, Collections.singletonList( diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AbstractSeriesAggregationScanOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AbstractSeriesAggregationScanOperator.java index d8ca5829403..dc8b5bf52d1 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AbstractSeriesAggregationScanOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AbstractSeriesAggregationScanOperator.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.source; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeId; @@ -57,7 +57,7 @@ public abstract class AbstractSeriesAggregationScanOperator extends AbstractData // We still think aggregator in SeriesAggregateScanOperator is a inputRaw step. // But in facing of statistics, it will invoke another method processStatistics() - protected final List<Aggregator> aggregators; + protected final List<TreeAggregator> aggregators; protected boolean finished = false; @@ -77,7 +77,7 @@ public abstract class AbstractSeriesAggregationScanOperator extends AbstractData OperatorContext context, SeriesScanUtil seriesScanUtil, int subSensorSize, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, boolean ascending, boolean outputEndTime, @@ -134,7 +134,7 @@ public abstract class AbstractSeriesAggregationScanOperator extends AbstractData // move to the next time window curTimeRange = timeRangeIterator.nextTimeRange(); // clear previous aggregation result - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { aggregator.reset(); } } @@ -230,7 +230,7 @@ public abstract class AbstractSeriesAggregationScanOperator extends AbstractData } protected void calcFromStatistics(Statistics timeStatistics, Statistics[] valueStatistics) { - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { if (aggregator.hasFinalResult()) { continue; } @@ -377,7 +377,7 @@ public abstract class AbstractSeriesAggregationScanOperator extends AbstractData if (outputEndTime) { dataTypes.add(TSDataType.INT64); } - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { dataTypes.addAll(Arrays.asList(aggregator.getOutputType())); } return dataTypes; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AlignedSeriesAggregationScanOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AlignedSeriesAggregationScanOperator.java index e513e77e6f1..67ebcfe3d20 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AlignedSeriesAggregationScanOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/AlignedSeriesAggregationScanOperator.java @@ -21,7 +21,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.source; import org.apache.iotdb.commons.path.AlignedFullPath; import org.apache.iotdb.db.queryengine.execution.MemoryEstimationHelper; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeId; @@ -47,7 +47,7 @@ public class AlignedSeriesAggregationScanOperator extends AbstractSeriesAggregat Ordering scanOrder, SeriesScanOptions scanOptions, OperatorContext context, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, GroupByTimeParameter groupByTimeParameter, long maxReturnSize, @@ -73,7 +73,7 @@ public class AlignedSeriesAggregationScanOperator extends AbstractSeriesAggregat boolean outputEndTime, SeriesScanOptions scanOptions, OperatorContext context, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, GroupByTimeParameter groupByTimeParameter, long maxReturnSize, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/SeriesAggregationScanOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/SeriesAggregationScanOperator.java index 72a67e0fe2b..3ba72b68fc0 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/SeriesAggregationScanOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/SeriesAggregationScanOperator.java @@ -21,7 +21,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.source; import org.apache.iotdb.commons.path.IFullPath; import org.apache.iotdb.db.queryengine.execution.MemoryEstimationHelper; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeId; @@ -54,7 +54,7 @@ public class SeriesAggregationScanOperator extends AbstractSeriesAggregationScan Ordering scanOrder, SeriesScanOptions scanOptions, OperatorContext context, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, GroupByTimeParameter groupByTimeParameter, long maxReturnSize, @@ -80,7 +80,7 @@ public class SeriesAggregationScanOperator extends AbstractSeriesAggregationScan boolean outputEndTime, SeriesScanOptions scanOptions, OperatorContext context, - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, ITimeRangeIterator timeRangeIterator, GroupByTimeParameter groupByTimeParameter, long maxReturnSize, 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 4a0b2873ffa..b3e073941f7 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 @@ -20,11 +20,14 @@ package org.apache.iotdb.db.queryengine.execution.operator.source.relational; import org.apache.iotdb.commons.path.AlignedFullPath; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; 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.AlignedSeriesScanUtil; -import org.apache.iotdb.db.queryengine.execution.operator.source.relational.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.operator.source.relational.aggregation.TableAggregator; +import org.apache.iotdb.db.queryengine.execution.operator.window.IWindow; +import org.apache.iotdb.db.queryengine.execution.operator.window.TimeWindow; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.queryengine.plan.planner.plan.parameter.GroupByTimeParameter; import org.apache.iotdb.db.queryengine.plan.planner.plan.parameter.SeriesScanOptions; @@ -34,10 +37,15 @@ import org.apache.iotdb.db.queryengine.plan.statement.component.Ordering; import org.apache.iotdb.db.storageengine.dataregion.read.IQueryDataSource; 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.statistics.Statistics; +import org.apache.tsfile.read.common.TimeRange; import org.apache.tsfile.read.common.block.TsBlock; import org.apache.tsfile.read.common.block.TsBlockBuilder; import org.apache.tsfile.read.common.block.column.LongColumn; +import org.apache.tsfile.utils.Pair; import org.apache.tsfile.write.schema.IMeasurementSchema; import java.io.IOException; @@ -46,11 +54,16 @@ import java.util.Optional; import java.util.stream.Collectors; import static org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil.appendAggregationResult; +import static org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil.calculateAggregationFromRawData; +import static org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil.isAllAggregatorsHasFinalResult; +import static org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil.process; +import static org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil.satisfiedTimeRange; 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 { - List<Aggregator> aggregators; + private final List<TableAggregator> aggregators; public static final LongColumn TIME_COLUMN_TEMPLATE = new LongColumn(1, Optional.empty(), new long[] {0}); @@ -98,7 +111,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation List<IMeasurementSchema> measurementSchemas, int maxTsBlockLineNum, int subSensorSize, - List<Aggregator> aggregators, + List<TableAggregator> aggregators, ITimeRangeIterator timeRangeIterator, boolean ascending, GroupByTimeParameter groupByTimeParameter, @@ -110,7 +123,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation context, null, subSensorSize, - aggregators, + null, timeRangeIterator, ascending, false, @@ -118,6 +131,8 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation maxReturnSize, canUseStatistics); + this.aggregators = aggregators; + this.sourceId = sourceId; this.operatorContext = context; this.columnSchemas = columnSchemas; @@ -159,7 +174,7 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation // move to the next time window curTimeRange = timeRangeIterator.nextTimeRange(); // clear previous aggregation result - for (Aggregator aggregator : aggregators) { + for (TableAggregator aggregator : aggregators) { aggregator.reset(); } } @@ -267,7 +282,267 @@ public class TableAggregationTableScanOperator extends AbstractSeriesAggregation @Override protected void updateResultTsBlock() { - appendAggregationResult( - resultTsBlockBuilder, aggregators, timeRangeIterator.currentOutputTime()); + appendAggregationResult(resultTsBlockBuilder, aggregators); + } + + protected boolean calcFromCachedData() { + return calcFromRawData(inputTsBlock); + } + + private boolean calcFromRawData(TsBlock tsBlock) { + Pair<Boolean, TsBlock> calcResult = + calculateAggregationFromRawData(tsBlock, aggregators, curTimeRange, ascending); + inputTsBlock = calcResult.getRight(); + return calcResult.getLeft(); + } + + /** + * Calculate aggregation value on the time range from the tsBlock containing raw data. + * + * @return left - whether the aggregation calculation of the current time range has done; right - + * remaining tsBlock + */ + public static Pair<Boolean, TsBlock> calculateAggregationFromRawData( + TsBlock inputTsBlock, + List<TableAggregator> aggregators, + TimeRange curTimeRange, + boolean ascending) { + if (inputTsBlock == null || inputTsBlock.isEmpty()) { + return new Pair<>(false, inputTsBlock); + } + + // check if the tsBlock does not contain points in current interval + if (satisfiedTimeRange(inputTsBlock, curTimeRange, ascending)) { + // skip points that cannot be calculated + if ((ascending && inputTsBlock.getStartTime() < curTimeRange.getMin()) + || (!ascending && inputTsBlock.getStartTime() > curTimeRange.getMax())) { + inputTsBlock = skipPointsOutOfTimeRange(inputTsBlock, curTimeRange, ascending); + } + + inputTsBlock = process(inputTsBlock, curTimeRange, aggregators); + } + + // judge whether the calculation finished + boolean isTsBlockOutOfBound = + inputTsBlock != null + && (ascending + ? inputTsBlock.getEndTime() > curTimeRange.getMax() + : inputTsBlock.getEndTime() < curTimeRange.getMin()); + return new Pair<>( + isAllAggregatorsHasFinalResult(aggregators) || isTsBlockOutOfBound, inputTsBlock); + } + + private static TsBlock process( + TsBlock inputTsBlock, TimeRange curTimeRange, List<TableAggregator> aggregators) { + // Get the row which need to be processed by aggregator + IWindow curWindow = new TimeWindow(curTimeRange); + Column timeColumn = inputTsBlock.getTimeColumn(); + int lastIndexToProcess = 0; + for (int i = 0; i < inputTsBlock.getPositionCount(); i++) { + if (!curWindow.satisfy(timeColumn, i)) { + break; + } + lastIndexToProcess = i; + } + + TsBlock inputRegion = inputTsBlock.getRegion(0, lastIndexToProcess + 1); + for (TableAggregator aggregator : aggregators) { + // current agg method has been calculated + if (aggregator.hasFinalResult()) { + continue; + } + aggregator.processBlock(inputRegion); + } + int lastReadRowIndex = lastIndexToProcess + 1; + if (lastReadRowIndex >= inputTsBlock.getPositionCount()) { + return null; + } else { + return inputTsBlock.subTsBlock(lastReadRowIndex); + } + } + + public static boolean isAllAggregatorsHasFinalResult(List<TableAggregator> aggregators) { + for (TableAggregator aggregator : aggregators) { + if (!aggregator.hasFinalResult()) { + return false; + } + } + return true; + } + + protected void calcFromStatistics(Statistics timeStatistics, Statistics[] valueStatistics) { + for (TreeAggregator aggregator : aggregators) { + if (aggregator.hasFinalResult()) { + continue; + } + aggregator.processStatistics(timeStatistics, valueStatistics); + } + } + + @SuppressWarnings({"squid:S3776", "squid:S135", "squid:S3740"}) + protected boolean readAndCalcFromFile() throws IOException { + // start stopwatch + long start = System.nanoTime(); + while (System.nanoTime() - start < leftRuntimeOfOneNextCall && seriesScanUtil.hasNextFile()) { + if (canUseStatistics && seriesScanUtil.canUseCurrentFileStatistics()) { + Statistics fileTimeStatistics = seriesScanUtil.currentFileTimeStatistics(); + if (fileTimeStatistics.getStartTime() > curTimeRange.getMax()) { + if (ascending) { + return true; + } else { + seriesScanUtil.skipCurrentFile(); + continue; + } + } + // calc from fileMetaData + if (curTimeRange.contains( + fileTimeStatistics.getStartTime(), fileTimeStatistics.getEndTime())) { + Statistics[] statisticsList = new Statistics[subSensorSize]; + for (int i = 0; i < subSensorSize; i++) { + statisticsList[i] = seriesScanUtil.currentFileStatistics(i); + } + calcFromStatistics(fileTimeStatistics, statisticsList); + seriesScanUtil.skipCurrentFile(); + if (isAllAggregatorsHasFinalResult(aggregators) && !isGroupByQuery) { + return true; + } else { + continue; + } + } + } + + // read chunk + if (readAndCalcFromChunk()) { + return true; + } + } + + return false; + } + + @SuppressWarnings({"squid:S3776", "squid:S135", "squid:S3740"}) + protected boolean readAndCalcFromChunk() throws IOException { + // start stopwatch + long start = System.nanoTime(); + while (System.nanoTime() - start < leftRuntimeOfOneNextCall && seriesScanUtil.hasNextChunk()) { + if (canUseStatistics && seriesScanUtil.canUseCurrentChunkStatistics()) { + Statistics chunkTimeStatistics = seriesScanUtil.currentChunkTimeStatistics(); + if (chunkTimeStatistics.getStartTime() > curTimeRange.getMax()) { + if (ascending) { + return true; + } else { + seriesScanUtil.skipCurrentChunk(); + continue; + } + } + // calc from chunkMetaData + if (curTimeRange.contains( + chunkTimeStatistics.getStartTime(), chunkTimeStatistics.getEndTime())) { + // calc from chunkMetaData + Statistics[] statisticsList = new Statistics[subSensorSize]; + for (int i = 0; i < subSensorSize; i++) { + statisticsList[i] = seriesScanUtil.currentChunkStatistics(i); + } + calcFromStatistics(chunkTimeStatistics, statisticsList); + seriesScanUtil.skipCurrentChunk(); + if (isAllAggregatorsHasFinalResult(aggregators) && !isGroupByQuery) { + return true; + } else { + continue; + } + } + } + + // read page + if (readAndCalcFromPage()) { + return true; + } + } + return false; + } + + @SuppressWarnings({"squid:S3776", "squid:S135", "squid:S3740"}) + protected boolean readAndCalcFromPage() throws IOException { + // start stopwatch + long start = System.nanoTime(); + try { + while (System.nanoTime() - start < leftRuntimeOfOneNextCall && seriesScanUtil.hasNextPage()) { + if (canUseStatistics && seriesScanUtil.canUseCurrentPageStatistics()) { + Statistics pageTimeStatistics = seriesScanUtil.currentPageTimeStatistics(); + // There is no more eligible points in current time range + if (pageTimeStatistics.getStartTime() > curTimeRange.getMax()) { + if (ascending) { + return true; + } else { + seriesScanUtil.skipCurrentPage(); + continue; + } + } + // can use pageHeader + if (curTimeRange.contains( + pageTimeStatistics.getStartTime(), pageTimeStatistics.getEndTime())) { + Statistics[] statisticsList = new Statistics[subSensorSize]; + for (int i = 0; i < subSensorSize; i++) { + statisticsList[i] = seriesScanUtil.currentPageStatistics(i); + } + calcFromStatistics(pageTimeStatistics, statisticsList); + seriesScanUtil.skipCurrentPage(); + if (isAllAggregatorsHasFinalResult(aggregators) && !isGroupByQuery) { + return true; + } else { + continue; + } + } + } + + // calc from page data + TsBlock tsBlock = seriesScanUtil.nextPage(); + if (tsBlock == null || tsBlock.isEmpty()) { + continue; + } + + // calc from raw data + if (calcFromRawData(tsBlock)) { + return true; + } + } + return false; + } finally { + leftRuntimeOfOneNextCall -= (System.nanoTime() - start); + } + } + + @Override + protected List<TSDataType> getResultDataTypes() { + // List<TSDataType> dataTypes = new ArrayList<>(); + // for (TableAggregator aggregator : aggregators) { + // dataTypes.add(aggregator.getType()); + // } + // return dataTypes; + return aggregators.stream().map(TableAggregator::getType).collect(Collectors.toList()); + } + + /** Append a row of aggregation results to the result tsBlock. */ + public static void appendAggregationResult( + TsBlockBuilder tsBlockBuilder, List<? extends TableAggregator> aggregators) { + // TimeColumnBuilder timeColumnBuilder = tsBlockBuilder.getTimeColumnBuilder(); + // Use start time of current time range as time column + // timeColumnBuilder.writeLong(outputTime); + ColumnBuilder[] columnBuilders = tsBlockBuilder.getValueColumnBuilders(); + int columnIndex = 0; + // if (endTime != INVALID_END_TIME) { + // columnBuilders[columnIndex].writeLong(endTime); + // columnIndex++; + // } + for (int i = 0; i < aggregators.size(); i++) { + // TableAggregator aggregator = aggregators.get(i); + // ColumnBuilder[] columnBuilder = new ColumnBuilder[aggregator.get, length]; + // columnBuilder[0] = columnBuilders[columnIndex++]; + // if (columnBuilder.length > 1) { + // columnBuilder[1] = columnBuilders[columnIndex++]; + // } + aggregators.get(i).evaluate(columnBuilders[i]); + } + tsBlockBuilder.declarePosition(); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/AggregationOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/AggregationOperator.java index 425f05f7867..3da8ba5c84f 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/AggregationOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/AggregationOperator.java @@ -47,7 +47,7 @@ public class AggregationOperator implements ProcessOperator { private final Operator child; - private final List<Aggregator> aggregators; + private final List<TableAggregator> aggregators; private final TsBlockBuilder resultBuilder; @@ -61,13 +61,13 @@ public class AggregationOperator implements ProcessOperator { private boolean finished = false; public AggregationOperator( - OperatorContext operatorContext, Operator child, List<Aggregator> aggregators) { + OperatorContext operatorContext, Operator child, List<TableAggregator> aggregators) { this.operatorContext = operatorContext; this.child = child; this.aggregators = aggregators; this.resultBuilder = new TsBlockBuilder( - aggregators.stream().map(Aggregator::getType).collect(toImmutableList())); + aggregators.stream().map(TableAggregator::getType).collect(toImmutableList())); this.resultColumnsBuilder = resultBuilder.getValueColumnBuilders(); this.memoryReservationManager = operatorContext @@ -96,7 +96,7 @@ public class AggregationOperator implements ProcessOperator { return null; } - for (Aggregator aggregator : aggregators) { + for (TableAggregator aggregator : aggregators) { aggregator.processBlock(block); } @@ -151,7 +151,7 @@ public class AggregationOperator implements ProcessOperator { public long ramBytesUsed() { return INSTANCE_SIZE + MemoryEstimationHelper.getEstimatedSizeOfAccountableObject(child) - + aggregators.stream().mapToLong(Aggregator::getEstimatedSize).count() + + aggregators.stream().mapToLong(TableAggregator::getEstimatedSize).count() + MemoryEstimationHelper.getEstimatedSizeOfAccountableObject(operatorContext) + resultBuilder.getRetainedSizeInBytes(); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/Aggregator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/TableAggregator.java similarity index 97% rename from iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/Aggregator.java rename to iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/TableAggregator.java index cbd465efd8e..527dab06b77 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/Aggregator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/aggregation/TableAggregator.java @@ -27,14 +27,14 @@ import java.util.OptionalInt; import static com.google.common.base.Preconditions.checkArgument; import static java.util.Objects.requireNonNull; -public class Aggregator { +public class TableAggregator { private final Accumulator accumulator; private final AggregationNode.Step step; private final TSDataType outputType; private final int[] inputChannels; private final OptionalInt maskChannel; - public Aggregator( + public TableAggregator( Accumulator accumulator, AggregationNode.Step step, TSDataType outputType, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/ConditionWindowManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/ConditionWindowManager.java index 70d251ef461..2804c80535a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/ConditionWindowManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/ConditionWindowManager.java @@ -21,7 +21,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.window; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory.KeepEvaluator; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.tsfile.block.column.Column; import org.apache.tsfile.enums.TSDataType; @@ -162,7 +162,7 @@ public class ConditionWindowManager implements IWindowManager { } @Override - public TsBlockBuilder createResultTsBlockBuilder(List<Aggregator> aggregators) { + public TsBlockBuilder createResultTsBlockBuilder(List<TreeAggregator> aggregators) { List<TSDataType> dataTypes = getResultDataTypes(aggregators); // Judge whether we need output endTime column. if (conditionWindow.isOutputEndTime()) { @@ -173,7 +173,7 @@ public class ConditionWindowManager implements IWindowManager { @Override public void appendAggregationResult( - TsBlockBuilder resultTsBlockBuilder, List<Aggregator> aggregators) { + TsBlockBuilder resultTsBlockBuilder, List<TreeAggregator> aggregators) { if (!keepEvaluator.apply(conditionWindow.getKeep())) { return; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/CountWindowManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/CountWindowManager.java index 3b86e9b462c..b81dd5af1e5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/CountWindowManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/CountWindowManager.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.window; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.tsfile.block.column.Column; import org.apache.tsfile.enums.TSDataType; @@ -116,7 +116,7 @@ public class CountWindowManager implements IWindowManager { } @Override - public TsBlockBuilder createResultTsBlockBuilder(List<Aggregator> aggregators) { + public TsBlockBuilder createResultTsBlockBuilder(List<TreeAggregator> aggregators) { List<TSDataType> dataTypes = getResultDataTypes(aggregators); // Judge whether we need output endTime column. if (countWindow.isNeedOutputEndTime()) { @@ -127,7 +127,7 @@ public class CountWindowManager implements IWindowManager { @Override public void appendAggregationResult( - TsBlockBuilder resultTsBlockBuilder, List<Aggregator> aggregators) { + TsBlockBuilder resultTsBlockBuilder, List<TreeAggregator> aggregators) { if (countWindow.getLeftCount() != 0) { return; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/IWindowManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/IWindowManager.java index c66a3bb725c..a1f32b7097d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/IWindowManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/IWindowManager.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.window; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.tsfile.block.column.ColumnBuilder; import org.apache.tsfile.enums.TSDataType; @@ -98,9 +98,9 @@ public interface IWindowManager { * @param aggregators the list of aggregators * @return Aggregation result column type list. */ - default List<TSDataType> getResultDataTypes(List<Aggregator> aggregators) { + default List<TSDataType> getResultDataTypes(List<TreeAggregator> aggregators) { List<TSDataType> dataTypes = new ArrayList<>(); - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { dataTypes.addAll(Arrays.asList(aggregator.getOutputType())); } return dataTypes; @@ -115,7 +115,7 @@ public interface IWindowManager { * @param aggregators the list of aggregators * @return TsBlockBuilder of resultSet */ - TsBlockBuilder createResultTsBlockBuilder(List<Aggregator> aggregators); + TsBlockBuilder createResultTsBlockBuilder(List<TreeAggregator> aggregators); /** * Used to append a row of aggregation result into the resultSet. @@ -127,7 +127,8 @@ public interface IWindowManager { * @param resultTsBlockBuilder tsBlockBuilder for resultSet * @param aggregators the list of aggregators */ - void appendAggregationResult(TsBlockBuilder resultTsBlockBuilder, List<Aggregator> aggregators); + void appendAggregationResult( + TsBlockBuilder resultTsBlockBuilder, List<TreeAggregator> aggregators); /** * Especially for TimeWindow, if there are no points belong to last TimeWindow, the last @@ -160,7 +161,7 @@ public interface IWindowManager { * @param endTime if the window doesn't need to output endTime, just assign -1 to endTime. */ default void outputAggregators( - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, TsBlockBuilder resultTsBlockBuilder, long startTime, long endTime) { @@ -173,7 +174,7 @@ public interface IWindowManager { columnBuilders[0].writeLong(endTime); columnIndex = 1; } - for (Aggregator aggregator : aggregators) { + for (TreeAggregator aggregator : aggregators) { ColumnBuilder[] columnBuilder = new ColumnBuilder[aggregator.getOutputType().length]; columnBuilder[0] = columnBuilders[columnIndex++]; if (columnBuilder.length > 1) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/SessionWindowManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/SessionWindowManager.java index cae70428af2..42ca3becedd 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/SessionWindowManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/SessionWindowManager.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.window; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.tsfile.block.column.Column; import org.apache.tsfile.enums.TSDataType; @@ -126,7 +126,7 @@ public class SessionWindowManager implements IWindowManager { } @Override - public TsBlockBuilder createResultTsBlockBuilder(List<Aggregator> aggregators) { + public TsBlockBuilder createResultTsBlockBuilder(List<TreeAggregator> aggregators) { List<TSDataType> dataTypes = getResultDataTypes(aggregators); // Judge whether we need output endTime column. if (isNeedOutputEndTime) { @@ -137,7 +137,7 @@ public class SessionWindowManager implements IWindowManager { @Override public void appendAggregationResult( - TsBlockBuilder resultTsBlockBuilder, List<Aggregator> aggregators) { + TsBlockBuilder resultTsBlockBuilder, List<TreeAggregator> aggregators) { long endTime = isNeedOutputEndTime ? sessionWindow.getEndTime() : -1; outputAggregators(aggregators, resultTsBlockBuilder, sessionWindow.getStartTime(), endTime); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/TimeWindowManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/TimeWindowManager.java index f5558e74fd4..9ecb42aec1f 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/TimeWindowManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/TimeWindowManager.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.window; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.operator.AggregationUtil; @@ -133,7 +133,7 @@ public class TimeWindowManager implements IWindowManager { } @Override - public TsBlockBuilder createResultTsBlockBuilder(List<Aggregator> aggregators) { + public TsBlockBuilder createResultTsBlockBuilder(List<TreeAggregator> aggregators) { List<TSDataType> dataTypes = getResultDataTypes(aggregators); // Judge whether we need output endTime column. if (this.needOutputEndTime) { @@ -144,7 +144,7 @@ public class TimeWindowManager implements IWindowManager { @Override public void appendAggregationResult( - TsBlockBuilder resultTsBlockBuilder, List<Aggregator> aggregators) { + TsBlockBuilder resultTsBlockBuilder, List<TreeAggregator> aggregators) { outputAggregators( aggregators, resultTsBlockBuilder, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/VariationWindowManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/VariationWindowManager.java index a4e1fecbc37..5c9ef777379 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/VariationWindowManager.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/window/VariationWindowManager.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.queryengine.execution.operator.window; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.read.common.block.TsBlockBuilder; @@ -78,7 +78,7 @@ public abstract class VariationWindowManager implements IWindowManager { } @Override - public TsBlockBuilder createResultTsBlockBuilder(List<Aggregator> aggregators) { + public TsBlockBuilder createResultTsBlockBuilder(List<TreeAggregator> aggregators) { List<TSDataType> dataTypes = getResultDataTypes(aggregators); // Judge whether we need output endTime column. if (variationWindow.isOutputEndTime()) { @@ -88,7 +88,7 @@ public abstract class VariationWindowManager implements IWindowManager { } public void appendAggregationResult( - TsBlockBuilder resultTsBlockBuilder, List<Aggregator> aggregators) { + TsBlockBuilder resultTsBlockBuilder, List<TreeAggregator> aggregators) { long endTime = variationWindow.isOutputEndTime() ? variationWindow.getEndTime() : -1; outputAggregators(aggregators, resultTsBlockBuilder, variationWindow.getStartTime(), endTime); 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 cf0036ffbad..2aa1a8a2bf8 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 @@ -36,7 +36,7 @@ import org.apache.iotdb.db.queryengine.common.NodeRef; import org.apache.iotdb.db.queryengine.common.TimeseriesContext; import org.apache.iotdb.db.queryengine.execution.aggregation.Accumulator; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.slidingwindow.SlidingWindowAggregatorFactory; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.driver.DataDriverContext; @@ -616,11 +616,11 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP (NonAlignedFullPath) IFullPath.convertToIFullPath(node.getSeriesPath()); boolean ascending = node.getScanOrder() == ASC; List<AggregationDescriptor> aggregationDescriptors = node.getAggregationDescriptorList(); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); aggregationDescriptors.forEach( o -> aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createAccumulator( o.getAggregationFuncName(), o.getAggregationType(), @@ -732,7 +732,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP LocalExecutionPlanContext context) { AlignedFullPath seriesPath = (AlignedFullPath) IFullPath.convertToIFullPath(alignedPath); boolean ascending = scanOrder == ASC; - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); boolean canUseStatistics = true; for (AggregationDescriptor descriptor : aggregationDescriptorList) { checkArgument( @@ -751,7 +751,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP canUseStatistics = false; } aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createAccumulator( descriptor.getAggregationFuncName(), descriptor.getAggregationType(), @@ -765,7 +765,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP new InputLocation[] {new InputLocation(0, seriesIndex)}))); } else if (expression instanceof TimestampOperand) { aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createAccumulator( descriptor.getAggregationFuncName(), descriptor.getAggregationType(), @@ -1790,7 +1790,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP "GroupByLevel descriptorList cannot be empty"); List<Operator> children = dealWithConsumeAllChildrenPipelineBreaker(node, context); boolean ascending = node.getScanOrder() == ASC; - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); Map<String, List<InputLocation>> layout = makeLayout(node); List<CrossSeriesAggregationDescriptor> aggregationDescriptors = node.getGroupByLevelDescriptors(); @@ -1807,7 +1807,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP descriptor.getInputExpressions().get(x).getExpressionString())) .collect(Collectors.toList()); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createAccumulator( descriptor.getAggregationFuncName(), descriptor.getAggregationType(), @@ -1850,12 +1850,12 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP boolean ascending = node.getScanOrder() == ASC; Map<String, List<InputLocation>> layout = makeLayout(node); List<List<String>> groups = new ArrayList<>(); - List<List<Aggregator>> groupedAggregators = new ArrayList<>(); + List<List<TreeAggregator>> groupedAggregators = new ArrayList<>(); int aggregatorCount = 0; for (Map.Entry<List<String>, List<CrossSeriesAggregationDescriptor>> entry : node.getTagValuesToAggregationDescriptors().entrySet()) { groups.add(entry.getKey()); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (CrossSeriesAggregationDescriptor aggregationDescriptor : entry.getValue()) { if (aggregationDescriptor == null) { aggregators.add(null); @@ -1867,7 +1867,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP .map(x -> context.getTypeProvider().getTreeModelType(x.getExpressionString())) .collect(Collectors.toList()); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createAccumulator( aggregationDescriptor.getAggregationFuncName(), aggregationDescriptor.getAggregationType(), @@ -1921,7 +1921,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP SlidingWindowAggregationOperator.class.getSimpleName()); Operator child = node.getChild().accept(this, context); boolean ascending = node.getScanOrder() == ASC; - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); Map<String, List<InputLocation>> layout = makeLayout(node); List<AggregationDescriptor> aggregationDescriptors = node.getAggregationDescriptorList(); for (AggregationDescriptor descriptor : aggregationDescriptors) { @@ -2019,12 +2019,12 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP "Aggregation descriptorList cannot be empty"); Operator child = node.getChild().accept(this, context); boolean ascending = node.getScanOrder() == ASC; - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); List<AggregationDescriptor> aggregationDescriptors = node.getAggregationDescriptorList(); for (AggregationDescriptor descriptor : node.getAggregationDescriptorList()) { List<InputLocation[]> inputLocationList = calcInputLocationList(descriptor, layout); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createAccumulator( descriptor.getAggregationFuncName(), descriptor.getAggregationType(), @@ -2144,13 +2144,13 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP "Aggregation descriptorList cannot be empty"); List<Operator> children = dealWithConsumeAllChildrenPipelineBreaker(node, context); boolean ascending = node.getScanOrder() == ASC; - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); Map<String, List<InputLocation>> layout = makeLayout(node); List<AggregationDescriptor> aggregationDescriptors = node.getAggregationDescriptorList(); for (AggregationDescriptor descriptor : node.getAggregationDescriptorList()) { List<InputLocation[]> inputLocationList = calcInputLocationList(descriptor, layout); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createAccumulator( descriptor.getAggregationFuncName(), descriptor.getAggregationType(), @@ -2923,7 +2923,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP SeriesAggregationScanOperator.class.getSimpleName()); // last_time, last_value - List<Aggregator> aggregators = LastQueryUtil.createAggregators(seriesPath.getSeriesType()); + List<TreeAggregator> aggregators = LastQueryUtil.createAggregators(seriesPath.getSeriesType()); ITimeRangeIterator timeRangeIterator = initTimeRangeIterator(null, false, false); long maxReturnSize = calculateMaxAggregationResultSizeForLastQuery(aggregators); @@ -2955,7 +2955,7 @@ public class OperatorTreeGenerator extends PlanVisitor<Operator, LocalExecutionP AlignedFullPath unCachedPath, LocalExecutionPlanContext context) { // last_time, last_value - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); boolean canUseStatistics = true; for (int i = 0; i < unCachedPath.getMeasurementList().size(); i++) { TSDataType dataType = unCachedPath.getSchemaList().get(i).getType(); 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 2f0152f160a..74c788b325b 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 @@ -51,13 +51,12 @@ import org.apache.iotdb.db.queryengine.execution.operator.schema.source.SchemaSo import org.apache.iotdb.db.queryengine.execution.operator.sink.IdentitySinkOperator; import org.apache.iotdb.db.queryengine.execution.operator.source.AlignedSeriesScanOperator; import org.apache.iotdb.db.queryengine.execution.operator.source.ExchangeOperator; -import org.apache.iotdb.db.queryengine.execution.operator.source.relational.TableAggregationTableScanOperator; import org.apache.iotdb.db.queryengine.execution.operator.source.relational.TableFullOuterJoinOperator; import org.apache.iotdb.db.queryengine.execution.operator.source.relational.TableInnerJoinOperator; import org.apache.iotdb.db.queryengine.execution.operator.source.relational.TableScanOperator; import org.apache.iotdb.db.queryengine.execution.operator.source.relational.aggregation.Accumulator; import org.apache.iotdb.db.queryengine.execution.operator.source.relational.aggregation.AggregationOperator; -import org.apache.iotdb.db.queryengine.execution.operator.source.relational.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.operator.source.relational.aggregation.TableAggregator; import org.apache.iotdb.db.queryengine.execution.relational.ColumnTransformerBuilder; import org.apache.iotdb.db.queryengine.plan.analyze.TypeProvider; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNode; @@ -946,7 +945,7 @@ public class TableOperatorGenerator extends PlanVisitor<Operator, LocalExecution AggregationNode node, Operator child, TypeProvider typeProvider, OperatorContext context) { Map<Symbol, AggregationNode.Aggregation> aggregationMap = node.getAggregations(); - ImmutableList.Builder<Aggregator> aggregatorBuilder = new ImmutableList.Builder<>(); + ImmutableList.Builder<TableAggregator> aggregatorBuilder = new ImmutableList.Builder<>(); Map<Symbol, Integer> childLayout = makeLayoutFromOutputSymbols(node.getChild().getOutputSymbols()); @@ -955,7 +954,11 @@ public class TableOperatorGenerator extends PlanVisitor<Operator, LocalExecution symbol -> aggregatorBuilder.add( buildAggregator( - childLayout, aggregationMap.get(symbol), node.getStep(), typeProvider))); + symbol, + childLayout, + aggregationMap.get(symbol), + node.getStep(), + typeProvider))); return new AggregationOperator(context, child, aggregatorBuilder.build()); } @@ -969,7 +972,8 @@ public class TableOperatorGenerator extends PlanVisitor<Operator, LocalExecution return outputMappings.buildOrThrow(); } - private Aggregator buildAggregator( + private TableAggregator buildAggregator( + Symbol aggregationSymbol, Map<Symbol, Integer> childLayout, AggregationNode.Aggregation aggregation, AggregationNode.Step step, @@ -999,10 +1003,10 @@ public class TableOperatorGenerator extends PlanVisitor<Operator, LocalExecution Collections.emptyMap(), true); - return new Aggregator( + return new TableAggregator( accumulator, step, - getTSDataType(aggregation.getResolvedFunction().getSignature().getReturnType()), + getTSDataType(typeProvider.getTableModelType(aggregationSymbol)), argumentChannels, OptionalInt.empty()); } @@ -1016,16 +1020,17 @@ public class TableOperatorGenerator extends PlanVisitor<Operator, LocalExecution .addOperatorContext( context.getNextOperatorId(), node.getPlanNodeId(), - TableAggregationTableScanOperator.class.getSimpleName()); + AggregationTableScanNode.class.getSimpleName()); - List<Aggregator> aggregators = new ArrayList<>(); + List<TableAggregator> aggregators = new ArrayList<>(); // TODO fix childLayout Map<Symbol, Integer> childLayout = new HashMap<>(); for (Map.Entry<Symbol, AggregationNode.Aggregation> entry : node.getAggregations().entrySet()) { - Aggregator aggregator = - buildAggregator(childLayout, entry.getValue(), node.getStep(), context.getTypeProvider()); + TableAggregator aggregator = + buildAggregator( + null, childLayout, entry.getValue(), node.getStep(), context.getTypeProvider()); aggregators.add(aggregator); } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationOperatorTest.java index f52eeb29b21..1993aee0a78 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AggregationOperatorTest.java @@ -29,7 +29,7 @@ import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.execution.aggregation.Accumulator; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -325,14 +325,14 @@ public class AggregationOperatorTest { new NonAlignedFullPath( IDeviceID.Factory.DEFAULT_FACTORY.create(AGGREGATION_OPERATOR_TEST_SG + ".device0"), new MeasurementSchema("sensor0", TSDataType.INT32)); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.PARTIAL))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.PARTIAL))); SeriesScanOptions.Builder scanOptionsBuilder = new SeriesScanOptions.Builder(); scanOptionsBuilder.withAllSensors(Collections.singleton("sensor0")); @@ -387,7 +387,7 @@ public class AggregationOperatorTest { children.add(seriesAggregationScanOperator1); children.add(seriesAggregationScanOperator2); - List<Aggregator> finalAggregators = new ArrayList<>(); + List<TreeAggregator> finalAggregators = new ArrayList<>(); List<Accumulator> accumulators = AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, @@ -397,7 +397,7 @@ public class AggregationOperatorTest { true); for (int i = 0; i < accumulators.size(); i++) { finalAggregators.add( - new Aggregator(accumulators.get(i), AggregationStep.FINAL, inputLocations.get(i))); + new TreeAggregator(accumulators.get(i), AggregationStep.FINAL, inputLocations.get(i))); } return new AggregationOperator( diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AlignedSeriesAggregationScanOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AlignedSeriesAggregationScanOperatorTest.java index 52adcccd0d7..2f38887e78e 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AlignedSeriesAggregationScanOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/AlignedSeriesAggregationScanOperatorTest.java @@ -28,7 +28,7 @@ import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -101,13 +101,13 @@ public class AlignedSeriesAggregationScanOperatorTest { @Test public void testAggregationWithoutTimeFilter() throws Exception { - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < measurementSchemas.size(); i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( TAggregationType.COUNT, Collections.singletonList(dataType), @@ -132,13 +132,13 @@ public class AlignedSeriesAggregationScanOperatorTest { @Test public void testAggregationWithoutTimeFilterOrderByTimeDesc() throws Exception { - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < measurementSchemas.size(); i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( TAggregationType.COUNT, Collections.singletonList(dataType), @@ -168,13 +168,13 @@ public class AlignedSeriesAggregationScanOperatorTest { aggregationTypes.add(TAggregationType.COUNT); aggregationTypes.add(TAggregationType.SUM); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < 2; i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( aggregationTypes.get(i), Collections.singletonList(dataType), @@ -206,13 +206,13 @@ public class AlignedSeriesAggregationScanOperatorTest { aggregationTypes.add(TAggregationType.MIN_TIME); aggregationTypes.add(TAggregationType.MAX_TIME); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < 6; i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( aggregationTypes.get(i), Collections.singletonList(dataType), @@ -248,13 +248,13 @@ public class AlignedSeriesAggregationScanOperatorTest { aggregationTypes.add(TAggregationType.MIN_TIME); aggregationTypes.add(TAggregationType.MAX_TIME); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < 6; i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( aggregationTypes.get(i), Collections.singletonList(dataType), @@ -282,13 +282,13 @@ public class AlignedSeriesAggregationScanOperatorTest { @Test public void testAggregationWithTimeFilter1() throws Exception { - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < measurementSchemas.size(); i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( TAggregationType.COUNT, Collections.singletonList(dataType), @@ -316,13 +316,13 @@ public class AlignedSeriesAggregationScanOperatorTest { @Test public void testAggregationWithTimeFilter2() throws Exception { Filter timeFilter = TimeFilterApi.ltEq(379); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < measurementSchemas.size(); i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( TAggregationType.COUNT, Collections.singletonList(dataType), @@ -349,13 +349,13 @@ public class AlignedSeriesAggregationScanOperatorTest { @Test public void testAggregationWithTimeFilter3() throws Exception { Filter timeFilter = FilterFactory.and(TimeFilterApi.gtEq(100), TimeFilterApi.ltEq(399)); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < measurementSchemas.size(); i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( TAggregationType.COUNT, Collections.singletonList(dataType), @@ -388,13 +388,13 @@ public class AlignedSeriesAggregationScanOperatorTest { aggregationTypes.add(TAggregationType.MIN_TIME); aggregationTypes.add(TAggregationType.MAX_TIME); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < 6; i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( aggregationTypes.get(i), Collections.singletonList(dataType), @@ -427,13 +427,13 @@ public class AlignedSeriesAggregationScanOperatorTest { int[] result = new int[] {100, 100, 100, 99}; GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < measurementSchemas.size(); i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( TAggregationType.COUNT, Collections.singletonList(dataType), @@ -467,13 +467,13 @@ public class AlignedSeriesAggregationScanOperatorTest { Filter timeFilter = FilterFactory.and(TimeFilterApi.gtEq(120), TimeFilterApi.ltEq(379)); GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); for (int i = 0; i < measurementSchemas.size(); i++) { TSDataType dataType = measurementSchemas.get(i).getType(); List<InputLocation[]> inputLocations = new ArrayList<>(); inputLocations.add(new InputLocation[] {new InputLocation(0, i)}); aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( TAggregationType.COUNT, Collections.singletonList(dataType), @@ -518,7 +518,7 @@ public class AlignedSeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); List<InputLocation[]> inputLocations = Collections.singletonList(new InputLocation[] {new InputLocation(0, 1)}); AccumulatorFactory.createBuiltinAccumulators( @@ -527,7 +527,8 @@ public class AlignedSeriesAggregationScanOperatorTest { Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE, inputLocations))); + .forEach( + o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE, inputLocations))); AlignedSeriesAggregationScanOperator seriesAggregationScanOperator = initAlignedSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -563,7 +564,7 @@ public class AlignedSeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); List<InputLocation[]> inputLocations = Collections.singletonList(new InputLocation[] {new InputLocation(0, 1)}); AccumulatorFactory.createBuiltinAccumulators( @@ -572,7 +573,8 @@ public class AlignedSeriesAggregationScanOperatorTest { Collections.emptyList(), Collections.emptyMap(), false) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE, inputLocations))); + .forEach( + o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE, inputLocations))); AlignedSeriesAggregationScanOperator seriesAggregationScanOperator = initAlignedSeriesAggregationScanOperator(aggregators, null, false, groupByTimeParameter); int count = 0; @@ -598,7 +600,7 @@ public class AlignedSeriesAggregationScanOperatorTest { new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 50), true); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); List<InputLocation[]> inputLocations = Collections.singletonList(new InputLocation[] {new InputLocation(0, 1)}); AccumulatorFactory.createBuiltinAccumulators( @@ -607,7 +609,8 @@ public class AlignedSeriesAggregationScanOperatorTest { Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE, inputLocations))); + .forEach( + o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE, inputLocations))); AlignedSeriesAggregationScanOperator seriesAggregationScanOperator = initAlignedSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -630,7 +633,7 @@ public class AlignedSeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 149, new TimeDuration(0, 50), new TimeDuration(0, 30), true); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); List<InputLocation[]> inputLocations = Collections.singletonList(new InputLocation[] {new InputLocation(0, 1)}); AccumulatorFactory.createBuiltinAccumulators( @@ -639,7 +642,8 @@ public class AlignedSeriesAggregationScanOperatorTest { Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE, inputLocations))); + .forEach( + o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE, inputLocations))); AlignedSeriesAggregationScanOperator seriesAggregationScanOperator = initAlignedSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -673,7 +677,7 @@ public class AlignedSeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 149, new TimeDuration(0, 50), new TimeDuration(0, 30), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); List<InputLocation[]> inputLocations = Collections.singletonList(new InputLocation[] {new InputLocation(0, 1)}); AccumulatorFactory.createBuiltinAccumulators( @@ -682,7 +686,8 @@ public class AlignedSeriesAggregationScanOperatorTest { Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE, inputLocations))); + .forEach( + o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE, inputLocations))); AlignedSeriesAggregationScanOperator seriesAggregationScanOperator = initAlignedSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -702,7 +707,7 @@ public class AlignedSeriesAggregationScanOperatorTest { } public AlignedSeriesAggregationScanOperator initAlignedSeriesAggregationScanOperator( - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, Filter timeFilter, boolean ascending, GroupByTimeParameter groupByTimeParameter) diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/HorizontallyConcatOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/HorizontallyConcatOperatorTest.java index a2cf16dbaef..498a7ddbdac 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/HorizontallyConcatOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/HorizontallyConcatOperatorTest.java @@ -27,7 +27,7 @@ import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -132,14 +132,14 @@ public class HorizontallyConcatOperatorTest { Arrays.asList(TAggregationType.COUNT, TAggregationType.SUM, TAggregationType.FIRST_VALUE); GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 10, new TimeDuration(0, 1), new TimeDuration(0, 1), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesScanOptions.Builder scanOptionsBuilder = new SeriesScanOptions.Builder(); scanOptionsBuilder.withAllSensors(allSensors); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryOperatorTest.java index 297fb4dca79..8a7711d1f4c 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryOperatorTest.java @@ -26,7 +26,7 @@ import org.apache.iotdb.commons.path.MeasurementPath; import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -96,10 +96,10 @@ public class LastQueryOperatorTest { @Test public void testLastQueryOperator1() throws Exception { try { - List<Aggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath1 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor0", TSDataType.INT32); - List<Aggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath2 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor1", TSDataType.INT32); Set<String> allSensors = Sets.newHashSet("sensor0", "sensor1"); @@ -224,10 +224,10 @@ public class LastQueryOperatorTest { @Test public void testLastQueryOperator2() { try { - List<Aggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath1 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor0", TSDataType.INT32); - List<Aggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath2 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor1", TSDataType.INT32); Set<String> allSensors = Sets.newHashSet("sensor0", "sensor1"); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryTreeSortOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryTreeSortOperatorTest.java index e7a187ce12f..29cd23bbb19 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryTreeSortOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/LastQueryTreeSortOperatorTest.java @@ -25,7 +25,7 @@ import org.apache.iotdb.commons.path.MeasurementPath; import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -97,10 +97,10 @@ public class LastQueryTreeSortOperatorTest { @Test public void testLastQuerySortOperatorAsc() { try { - List<Aggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath1 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor0", TSDataType.INT32); - List<Aggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath2 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor1", TSDataType.INT32); Set<String> allSensors = Sets.newHashSet("sensor0", "sensor1"); @@ -225,10 +225,10 @@ public class LastQueryTreeSortOperatorTest { @Test public void testLastQuerySortOperatorDesc() { try { - List<Aggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators1 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath1 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor0", TSDataType.INT32); - List<Aggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators2 = LastQueryUtil.createAggregators(TSDataType.INT32); MeasurementPath measurementPath2 = new MeasurementPath(SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor1", TSDataType.INT32); Set<String> allSensors = Sets.newHashSet("sensor0", "sensor1"); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorMemoryTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorMemoryTest.java index e22a2df6a0a..6397421bbd0 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorMemoryTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorMemoryTest.java @@ -30,7 +30,7 @@ import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.timerangeiterator.ITimeRangeIterator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; @@ -1193,11 +1193,11 @@ public class OperatorMemoryTest { PlanNodeId planNodeId = new PlanNodeId("1"); driverContext.addOperatorContext(1, planNodeId, SeriesScanOperator.class.getSimpleName()); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); aggregationDescriptors.forEach( o -> aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( o.getAggregationType(), Collections.singletonList(measurementPath.getSeriesType()), @@ -1249,11 +1249,11 @@ public class OperatorMemoryTest { AggregationStep.FINAL, Collections.singletonList(new TimeSeriesOperand(measurementPath)))); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); aggregationDescriptors.forEach( o -> aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( o.getAggregationType(), Collections.singletonList(measurementPath.getSeriesType()), @@ -1322,11 +1322,11 @@ public class OperatorMemoryTest { AggregationStep.FINAL, Collections.singletonList(new TimeSeriesOperand(measurementPath)))); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); aggregationDescriptors.forEach( o -> aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( o.getAggregationType(), Collections.singletonList(measurementPath.getSeriesType()), @@ -1402,11 +1402,11 @@ public class OperatorMemoryTest { AggregationStep.FINAL, Collections.singletonList(new TimeSeriesOperand(measurementPath)))); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); aggregationDescriptors.forEach( o -> aggregators.add( - new Aggregator( + new TreeAggregator( AccumulatorFactory.createBuiltinAccumulator( o.getAggregationType(), Collections.singletonList(measurementPath.getSeriesType()), diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/RawDataAggregationOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/RawDataAggregationOperatorTest.java index 163d2692e40..0410265bcea 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/RawDataAggregationOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/RawDataAggregationOperatorTest.java @@ -30,7 +30,7 @@ import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.execution.aggregation.Accumulator; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -980,7 +980,7 @@ public class RawDataAggregationOperatorTest { new SingleColumnMerger(new InputLocation(1, 0), new AscTimeComparator())), new AscTimeComparator()); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); List<Accumulator> accumulators = AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, @@ -990,7 +990,7 @@ public class RawDataAggregationOperatorTest { true); for (int i = 0; i < accumulators.size(); i++) { aggregators.add( - new Aggregator(accumulators.get(i), AggregationStep.SINGLE, inputLocations.get(i))); + new TreeAggregator(accumulators.get(i), AggregationStep.SINGLE, inputLocations.get(i))); } return new RawDataAggregationOperator( driverContext.getOperatorContexts().get(3), diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SeriesAggregationScanOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SeriesAggregationScanOperatorTest.java index 343b53490db..42ddfcf64bb 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SeriesAggregationScanOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SeriesAggregationScanOperatorTest.java @@ -27,7 +27,7 @@ import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -99,14 +99,14 @@ public class SeriesAggregationScanOperatorTest { @Test public void testAggregationWithoutTimeFilter() throws Exception { List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, null); int count = 0; @@ -123,14 +123,14 @@ public class SeriesAggregationScanOperatorTest { @Test public void testAggregationWithoutTimeFilterOrderByTimeDesc() throws Exception { List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, false, null); int count = 0; @@ -149,14 +149,14 @@ public class SeriesAggregationScanOperatorTest { List<TAggregationType> aggregationTypes = new ArrayList<>(); aggregationTypes.add(TAggregationType.COUNT); aggregationTypes.add(TAggregationType.SUM); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, null); int count = 0; @@ -180,14 +180,14 @@ public class SeriesAggregationScanOperatorTest { aggregationTypes.add(TAggregationType.MAX_TIME); aggregationTypes.add(TAggregationType.MAX_VALUE); aggregationTypes.add(TAggregationType.MIN_VALUE); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, null); int count = 0; @@ -216,14 +216,14 @@ public class SeriesAggregationScanOperatorTest { aggregationTypes.add(TAggregationType.MAX_VALUE); aggregationTypes.add(TAggregationType.MIN_VALUE); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), false) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, false, null); int count = 0; @@ -246,14 +246,14 @@ public class SeriesAggregationScanOperatorTest { public void testAggregationWithTimeFilter1() throws Exception { List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); Filter timeFilter = TimeFilterApi.gtEq(120); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, timeFilter, true, null); @@ -273,14 +273,14 @@ public class SeriesAggregationScanOperatorTest { Filter timeFilter = TimeFilterApi.ltEq(379); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, timeFilter, true, null); int count = 0; @@ -299,14 +299,14 @@ public class SeriesAggregationScanOperatorTest { Filter timeFilter = FilterFactory.and(TimeFilterApi.gtEq(100), TimeFilterApi.ltEq(399)); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, timeFilter, true, null); int count = 0; @@ -330,14 +330,14 @@ public class SeriesAggregationScanOperatorTest { aggregationTypes.add(TAggregationType.MAX_VALUE); aggregationTypes.add(TAggregationType.MIN_VALUE); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); Filter timeFilter = FilterFactory.and(TimeFilterApi.gtEq(100), TimeFilterApi.ltEq(399)); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, timeFilter, true, null); @@ -363,14 +363,14 @@ public class SeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -396,14 +396,14 @@ public class SeriesAggregationScanOperatorTest { new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, timeFilter, true, groupByTimeParameter); int count = 0; @@ -438,14 +438,14 @@ public class SeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -483,14 +483,14 @@ public class SeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 100), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), false) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, false, groupByTimeParameter); int count = 0; @@ -518,14 +518,14 @@ public class SeriesAggregationScanOperatorTest { new GroupByTimeParameter(0, 399, new TimeDuration(0, 100), new TimeDuration(0, 50), true); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -551,14 +551,14 @@ public class SeriesAggregationScanOperatorTest { new GroupByTimeParameter(0, 149, new TimeDuration(0, 50), new TimeDuration(0, 30), true); List<TAggregationType> aggregationTypes = Collections.singletonList(TAggregationType.COUNT); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -594,14 +594,14 @@ public class SeriesAggregationScanOperatorTest { GroupByTimeParameter groupByTimeParameter = new GroupByTimeParameter(0, 149, new TimeDuration(0, 50), new TimeDuration(0, 30), true); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( aggregationTypes, TSDataType.INT32, Collections.emptyList(), Collections.emptyMap(), true) - .forEach(o -> aggregators.add(new Aggregator(o, AggregationStep.SINGLE))); + .forEach(o -> aggregators.add(new TreeAggregator(o, AggregationStep.SINGLE))); SeriesAggregationScanOperator seriesAggregationScanOperator = initSeriesAggregationScanOperator(aggregators, null, true, groupByTimeParameter); int count = 0; @@ -623,7 +623,7 @@ public class SeriesAggregationScanOperatorTest { } public SeriesAggregationScanOperator initSeriesAggregationScanOperator( - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, Filter timeFilter, boolean ascending, GroupByTimeParameter groupByTimeParameter) diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SlidingWindowAggregationOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SlidingWindowAggregationOperatorTest.java index 2bdcb28957e..ec18f86716c 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SlidingWindowAggregationOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/SlidingWindowAggregationOperatorTest.java @@ -28,7 +28,7 @@ import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.execution.aggregation.AccumulatorFactory; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.aggregation.slidingwindow.SlidingWindowAggregatorFactory; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; @@ -241,7 +241,7 @@ public class SlidingWindowAggregationOperatorTest { IDeviceID.Factory.DEFAULT_FACTORY.create(AGGREGATION_OPERATOR_TEST_SG + ".device0"), new MeasurementSchema("sensor0", TSDataType.INT32)); - List<Aggregator> aggregators = new ArrayList<>(); + List<TreeAggregator> aggregators = new ArrayList<>(); AccumulatorFactory.createBuiltinAccumulators( leafAggregationTypes, TSDataType.INT32, @@ -249,7 +249,8 @@ public class SlidingWindowAggregationOperatorTest { Collections.emptyMap(), ascending) .forEach( - accumulator -> aggregators.add(new Aggregator(accumulator, AggregationStep.PARTIAL))); + accumulator -> + aggregators.add(new TreeAggregator(accumulator, AggregationStep.PARTIAL))); SeriesScanOptions.Builder scanOptionsBuilder = new SeriesScanOptions.Builder(); scanOptionsBuilder.withAllSensors(Collections.singleton("sensor0")); @@ -268,7 +269,7 @@ public class SlidingWindowAggregationOperatorTest { seriesAggregationScanOperator.initQueryDataSource( new QueryDataSource(seqResources, unSeqResources)); - List<Aggregator> finalAggregators = new ArrayList<>(); + List<TreeAggregator> finalAggregators = new ArrayList<>(); for (int i = 0; i < rootAggregationTypes.size(); i++) { finalAggregators.add( SlidingWindowAggregatorFactory.createSlidingWindowAggregator( diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/UpdateLastCacheOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/UpdateLastCacheOperatorTest.java index d991caff54b..1cd55520985 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/UpdateLastCacheOperatorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/operator/UpdateLastCacheOperatorTest.java @@ -26,7 +26,7 @@ import org.apache.iotdb.commons.path.MeasurementPath; import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.common.QueryId; -import org.apache.iotdb.db.queryengine.execution.aggregation.Aggregator; +import org.apache.iotdb.db.queryengine.execution.aggregation.TreeAggregator; import org.apache.iotdb.db.queryengine.execution.driver.DriverContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceStateMachine; @@ -97,7 +97,7 @@ public class UpdateLastCacheOperatorTest { @Test public void testUpdateLastCacheOperatorTestWithoutTimeFilter() { try { - List<Aggregator> aggregators = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators = LastQueryUtil.createAggregators(TSDataType.INT32); UpdateLastCacheOperator updateLastCacheOperator = initUpdateLastCacheOperator(aggregators, null, false, null); @@ -126,7 +126,7 @@ public class UpdateLastCacheOperatorTest { @Test public void testUpdateLastCacheOperatorTestWithTimeFilter1() { try { - List<Aggregator> aggregators = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators = LastQueryUtil.createAggregators(TSDataType.INT32); Filter timeFilter = TimeFilterApi.gtEq(200); UpdateLastCacheOperator updateLastCacheOperator = initUpdateLastCacheOperator(aggregators, timeFilter, false, null); @@ -156,7 +156,7 @@ public class UpdateLastCacheOperatorTest { @Test public void testUpdateLastCacheOperatorTestWithTimeFilter2() { try { - List<Aggregator> aggregators = LastQueryUtil.createAggregators(TSDataType.INT32); + List<TreeAggregator> aggregators = LastQueryUtil.createAggregators(TSDataType.INT32); Filter timeFilter = TimeFilterApi.ltEq(120); UpdateLastCacheOperator updateLastCacheOperator = initUpdateLastCacheOperator(aggregators, timeFilter, false, null); @@ -184,7 +184,7 @@ public class UpdateLastCacheOperatorTest { } public UpdateLastCacheOperator initUpdateLastCacheOperator( - List<Aggregator> aggregators, + List<TreeAggregator> aggregators, Filter timeFilter, boolean ascending, GroupByTimeParameter groupByTimeParameter)
