This is an automated email from the ASF dual-hosted git repository. jackietien pushed a commit to branch IOTDB-3883 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 9a6f0bad93ef2880c160fd9488d0c4b1d8a95e43 Author: JackieTien97 <[email protected]> AuthorDate: Fri Jul 22 16:55:14 2022 +0800 [IOTDB-3883] Support Order by timeseries in last query --- .../operator/process/ProcessOperator.java | 4 +- ...Operator.java => LastQueryCollectOperator.java} | 35 ++-- .../process/last/LastQueryMergeOperator.java | 178 +++++++++++++++++++-- .../operator/process/last/LastQueryOperator.java | 22 +-- .../process/last/LastQuerySortOperator.java | 38 +++-- .../operator/{ => process/last}/LastQueryUtil.java | 21 ++- .../process/last/UpdateLastCacheOperator.java | 1 - .../db/mpp/plan/planner/LocalExecutionPlanner.java | 155 +++++++++++++----- .../db/mpp/plan/planner/LogicalPlanBuilder.java | 4 +- .../planner/distribution/ExchangeNodeAdder.java | 4 +- .../plan/planner/distribution/SourceRewriter.java | 12 +- .../mpp/plan/planner/plan/node/PlanNodeType.java | 20 ++- .../db/mpp/plan/planner/plan/node/PlanVisitor.java | 12 +- .../node/process/last/LastQueryCollectNode.java | 108 +++++++++++++ .../process/{ => last}/LastQueryMergeNode.java | 52 ++---- .../LastQueryNode.java} | 28 ++-- .../iotdb/db/query/executor/LastQueryExecutor.java | 2 +- .../operator/AggregationOperatorTest.java | 2 +- .../operator/LastCacheScanOperatorTest.java | 93 ----------- .../execution/operator/LastQueryOperatorTest.java | 28 ++-- ...torTest.java => LastQuerySortOperatorTest.java} | 62 ++++--- .../operator/UpdateLastCacheOperatorTest.java | 1 + .../db/mpp/plan/plan/QueryLogicalPlanUtil.java | 8 +- 23 files changed, 573 insertions(+), 317 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/ProcessOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/ProcessOperator.java index a6f64ff295..6205cb766d 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/ProcessOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/ProcessOperator.java @@ -21,6 +21,4 @@ package org.apache.iotdb.db.mpp.execution.operator.process; import org.apache.iotdb.db.mpp.execution.operator.Operator; // TODO should think about what interfaces should this ProcessOperator have -public interface ProcessOperator extends Operator { - -} +public interface ProcessOperator extends Operator {} diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryMergeOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryCollectOperator.java similarity index 73% copy from server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryMergeOperator.java copy to server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryCollectOperator.java index d6998ef1af..7f011c32cf 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryMergeOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryCollectOperator.java @@ -18,18 +18,17 @@ */ package org.apache.iotdb.db.mpp.execution.operator.process.last; -import com.google.common.util.concurrent.ListenableFuture; -import org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.Operator; import org.apache.iotdb.db.mpp.execution.operator.OperatorContext; import org.apache.iotdb.db.mpp.execution.operator.process.ProcessOperator; import org.apache.iotdb.tsfile.read.common.block.TsBlock; -import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder; + +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import java.util.List; -// merge all last query result from different data regions, it will select max time for the same time-series -public class LastQueryMergeOperator implements ProcessOperator { +public class LastQueryCollectOperator implements ProcessOperator { private final OperatorContext operatorContext; @@ -39,14 +38,11 @@ public class LastQueryMergeOperator implements ProcessOperator { private int currentIndex; - private TsBlockBuilder tsBlockBuilder; - - public LastQueryMergeOperator(OperatorContext operatorContext, List<Operator> children) { + public LastQueryCollectOperator(OperatorContext operatorContext, List<Operator> children) { this.operatorContext = operatorContext; this.children = children; this.inputOperatorsCount = children.size(); this.currentIndex = 0; - this.tsBlockBuilder = LastQueryUtil.createTsBlockBuilder(); } @Override @@ -56,26 +52,37 @@ public class LastQueryMergeOperator implements ProcessOperator { @Override public ListenableFuture<?> isBlocked() { - return ProcessOperator.super.isBlocked(); + if (currentIndex < inputOperatorsCount) { + return children.get(currentIndex).isBlocked(); + } else { + return Futures.immediateVoidFuture(); + } } @Override public TsBlock next() { - return null; + if (children.get(currentIndex).hasNext()) { + return children.get(currentIndex).next(); + } else { + currentIndex++; + return null; + } } @Override public boolean hasNext() { - return false; + return currentIndex < inputOperatorsCount; } @Override public void close() throws Exception { - ProcessOperator.super.close(); + for (Operator child : children) { + child.close(); + } } @Override public boolean isFinished() { - return false; + return !hasNext(); } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryMergeOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryMergeOperator.java index d6998ef1af..fadf0b4d5c 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryMergeOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryMergeOperator.java @@ -18,17 +18,26 @@ */ package org.apache.iotdb.db.mpp.execution.operator.process.last; -import com.google.common.util.concurrent.ListenableFuture; -import org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.Operator; import org.apache.iotdb.db.mpp.execution.operator.OperatorContext; import org.apache.iotdb.db.mpp.execution.operator.process.ProcessOperator; import org.apache.iotdb.tsfile.read.common.block.TsBlock; import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder; +import org.apache.iotdb.tsfile.utils.Binary; + +import com.google.common.util.concurrent.ListenableFuture; +import java.util.ArrayList; +import java.util.Comparator; import java.util.List; +import java.util.TreeMap; -// merge all last query result from different data regions, it will select max time for the same time-series +import static com.google.common.util.concurrent.Futures.successfulAsList; +import static org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil.appendLastValue; +import static org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil.getTimeSeries; + +// merge all last query result from different data regions, it will select max time for the same +// time-series public class LastQueryMergeOperator implements ProcessOperator { private final OperatorContext operatorContext; @@ -37,16 +46,38 @@ public class LastQueryMergeOperator implements ProcessOperator { private final int inputOperatorsCount; - private int currentIndex; - private TsBlockBuilder tsBlockBuilder; - public LastQueryMergeOperator(OperatorContext operatorContext, List<Operator> children) { + /** TsBlock from child operator. Only one cache now. */ + private final TsBlock[] inputTsBlocks; + + /** start index for each input TsBlocks and size of it is equal to inputTsBlocks */ + private final int[] inputIndex; + + /** + * Represent whether there are more tsBlocks from ith child operator. If all elements in + * noMoreTsBlocks[] are true and inputTsBlocks[] are consumed completely, this operator is + * finished. + */ + private final boolean[] noMoreTsBlocks; + + private boolean finished = false; + + private final Comparator<Binary> comparator; + + private final TreeMap<Binary, Location> timeSeriesSelector; + + public LastQueryMergeOperator( + OperatorContext operatorContext, List<Operator> children, Comparator<Binary> comparator) { this.operatorContext = operatorContext; this.children = children; this.inputOperatorsCount = children.size(); - this.currentIndex = 0; this.tsBlockBuilder = LastQueryUtil.createTsBlockBuilder(); + this.inputTsBlocks = new TsBlock[this.inputOperatorsCount]; + this.inputIndex = new int[this.inputOperatorsCount]; + this.noMoreTsBlocks = new boolean[this.inputOperatorsCount]; + this.comparator = comparator; + this.timeSeriesSelector = new TreeMap<>(comparator); } @Override @@ -56,26 +87,151 @@ public class LastQueryMergeOperator implements ProcessOperator { @Override public ListenableFuture<?> isBlocked() { - return ProcessOperator.super.isBlocked(); + List<ListenableFuture<?>> listenableFutures = new ArrayList<>(); + for (int i = 0; i < inputOperatorsCount; i++) { + if (!noMoreTsBlocks[i] && empty(i)) { + ListenableFuture<?> blocked = children.get(i).isBlocked(); + if (!blocked.isDone()) { + listenableFutures.add(blocked); + } + } + } + return listenableFutures.isEmpty() ? NOT_BLOCKED : successfulAsList(listenableFutures); } @Override public TsBlock next() { - return null; + + // end time series for returned TsBlock this time, it's the min/max end time series among all + // the children + // TsBlocks order by asc/desc + Binary currentEndTimeSeries = null; + boolean init = false; + // get TsBlock for each input, put their time series into TimeSeriesSelector and then use the + // min/max TimeSeries + // among all the input TsBlock as the current output TsBlock's endTimeSeries. + for (int i = 0; i < inputOperatorsCount; i++) { + if (!noMoreTsBlocks[i] && empty(i)) { + if (children.get(i).hasNext()) { + inputIndex[i] = 0; + inputTsBlocks[i] = children.get(i).next(); + if (!empty(i)) { + int rowSize = inputTsBlocks[i].getPositionCount(); + for (int row = 0; row < rowSize; row++) { + Binary key = getTimeSeries(inputTsBlocks[i], row); + Location location = timeSeriesSelector.get(key); + if (location == null + || inputTsBlocks[i].getTimeByIndex(row) + > inputTsBlocks[location.tsBlockIndex].getTimeByIndex(location.rowIndex)) { + timeSeriesSelector.put(key, new Location(i, row)); + } + } + } else { + // child operator has next but return an empty TsBlock which means that it may not + // finish calculation in given time slice. + // In such case, LastQueryMergeOperator can't go on calculating, so we just return null. + // We can also use the while loop here to continuously call the hasNext() and next() + // methods of the child operator until its hasNext() returns false or the next() gets + // the data that is not empty, but this will cause the execution time of the while loop + // to be uncontrollable and may exceed all allocated time slice + return null; + } + } else { // no more tsBlock + noMoreTsBlocks[i] = true; + inputTsBlocks[i] = null; + } + } + // update the currentEndTimeSeries if the TsBlock is not empty + if (!empty(i)) { + Binary endTimeSeries = + getTimeSeries(inputTsBlocks[i], inputTsBlocks[i].getPositionCount() - 1); + currentEndTimeSeries = + init + ? (comparator.compare(currentEndTimeSeries, endTimeSeries) < 0 + ? currentEndTimeSeries + : endTimeSeries) + : endTimeSeries; + init = true; + } + } + + if (timeSeriesSelector.isEmpty()) { + return tsBlockBuilder.build(); + } + + while (!timeSeriesSelector.isEmpty() + && (comparator.compare(timeSeriesSelector.firstKey(), currentEndTimeSeries) <= 0)) { + Location location = timeSeriesSelector.pollFirstEntry().getValue(); + appendLastValue(tsBlockBuilder, inputTsBlocks[location.tsBlockIndex], location.rowIndex); + tsBlockBuilder.declarePosition(); + } + + TsBlock res = tsBlockBuilder.build(); + tsBlockBuilder.reset(); + return res; } @Override public boolean hasNext() { + if (finished) { + return false; + } + for (int i = 0; i < inputOperatorsCount; i++) { + if (!empty(i)) { + return true; + } else if (!noMoreTsBlocks[i]) { + if (children.get(i).hasNext()) { + return true; + } else { + noMoreTsBlocks[i] = true; + inputTsBlocks[i] = null; + } + } + } return false; } @Override public void close() throws Exception { - ProcessOperator.super.close(); + for (Operator child : children) { + child.close(); + } + tsBlockBuilder = null; } @Override public boolean isFinished() { - return false; + if (finished) { + return true; + } + finished = true; + + for (int i = 0; i < inputOperatorsCount; i++) { + // has more tsBlock output from children[i] or has cached tsBlock in inputTsBlocks[i] + if (!noMoreTsBlocks[i] || !empty(i)) { + finished = false; + break; + } + } + return finished; + } + + /** + * If the tsBlock of columnIndex is null or has no more data in the tsBlock, return true; else + * return false; + */ + private boolean empty(int columnIndex) { + return inputTsBlocks[columnIndex] == null + || inputTsBlocks[columnIndex].getPositionCount() == inputIndex[columnIndex]; + } + + private static class Location { + int tsBlockIndex; + int rowIndex; + + public Location(int tsBlockIndex, int rowIndex) { + this.tsBlockIndex = tsBlockIndex; + this.rowIndex = rowIndex; + } } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryOperator.java index 7152347293..517ebc4032 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryOperator.java @@ -18,16 +18,15 @@ */ package org.apache.iotdb.db.mpp.execution.operator.process.last; -import org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.Operator; import org.apache.iotdb.db.mpp.execution.operator.OperatorContext; import org.apache.iotdb.db.mpp.execution.operator.process.ProcessOperator; import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor; import org.apache.iotdb.tsfile.read.common.block.TsBlock; +import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder; import java.util.ArrayList; import java.util.List; @@ -35,11 +34,11 @@ import java.util.concurrent.TimeUnit; import static com.google.common.util.concurrent.Futures.successfulAsList; - // collect all last query result in the same data region and there is no order guarantee public class LastQueryOperator implements ProcessOperator { - private static final int MAX_DETECT_COUNT = TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(); + private static final int MAX_DETECT_COUNT = + TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(); private final OperatorContext operatorContext; @@ -51,8 +50,10 @@ public class LastQueryOperator implements ProcessOperator { private TsBlockBuilder tsBlockBuilder; - - public LastQueryOperator(OperatorContext operatorContext, List<UpdateLastCacheOperator> children, TsBlockBuilder builder) { + public LastQueryOperator( + OperatorContext operatorContext, + List<UpdateLastCacheOperator> children, + TsBlockBuilder builder) { this.operatorContext = operatorContext; this.children = children; this.inputOperatorsCount = children.size(); @@ -85,7 +86,8 @@ public class LastQueryOperator implements ProcessOperator { @Override public TsBlock next() { - // we have consumed up data from children Operator, just return all remaining cached data in tsBlockBuilder + // we have consumed up data from children Operator, just return all remaining cached data in + // tsBlockBuilder if (currentIndex >= inputOperatorsCount) { return tsBlockBuilder.build(); } @@ -96,7 +98,9 @@ public class LastQueryOperator implements ProcessOperator { int endIndex = getEndIndex(); - while ((System.nanoTime() - start < maxRuntime) && (currentIndex < endIndex) && !tsBlockBuilder.isFull()) { + while ((System.nanoTime() - start < maxRuntime) + && (currentIndex < endIndex) + && !tsBlockBuilder.isFull()) { if (children.get(currentIndex).hasNext()) { TsBlock tsBlock = children.get(currentIndex).next(); if (tsBlock == null) { @@ -134,4 +138,4 @@ public class LastQueryOperator implements ProcessOperator { private int getEndIndex() { return currentIndex + Math.min(MAX_DETECT_COUNT, inputOperatorsCount - currentIndex); } -} \ No newline at end of file +} diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQuerySortOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQuerySortOperator.java index 458e2c70ef..2945856933 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQuerySortOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQuerySortOperator.java @@ -18,10 +18,6 @@ */ package org.apache.iotdb.db.mpp.execution.operator.process.last; - -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.Operator; import org.apache.iotdb.db.mpp.execution.operator.OperatorContext; import org.apache.iotdb.db.mpp.execution.operator.process.ProcessOperator; @@ -30,18 +26,23 @@ import org.apache.iotdb.tsfile.read.common.block.TsBlock; import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder; import org.apache.iotdb.tsfile.utils.Binary; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; + import java.util.ArrayList; import java.util.Comparator; import java.util.List; import java.util.concurrent.TimeUnit; import static com.google.common.util.concurrent.Futures.successfulAsList; -import static org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil.compareTimeSeries; +import static org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil.compareTimeSeries; -// collect all last query result in the same data region and sort them according to the time-series's alphabetical order +// collect all last query result in the same data region and sort them according to the +// time-series's alphabetical order public class LastQuerySortOperator implements ProcessOperator { - private static final int MAX_DETECT_COUNT = TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(); + private static final int MAX_DETECT_COUNT = + TSFileDescriptor.getInstance().getConfig().getMaxTsBlockLineNumber(); // we must make sure that data in cachedTsBlock has already been sorted // values that have last cache @@ -68,7 +69,11 @@ public class LastQuerySortOperator implements ProcessOperator { // used to cache previous TsBlock get from children private TsBlock previousTsBlock; - public LastQuerySortOperator(OperatorContext operatorContext, TsBlock cachedTsBlock, List<UpdateLastCacheOperator> children, Comparator<Binary> timeSeriesComparator) { + public LastQuerySortOperator( + OperatorContext operatorContext, + TsBlock cachedTsBlock, + List<UpdateLastCacheOperator> children, + Comparator<Binary> timeSeriesComparator) { this.cachedTsBlock = cachedTsBlock; this.cachedTsBlockSize = cachedTsBlock.getPositionCount(); this.operatorContext = operatorContext; @@ -104,7 +109,8 @@ public class LastQuerySortOperator implements ProcessOperator { @Override public TsBlock next() { - // we have consumed up data from children Operator, just return all remaining cached data in cachedTsBlock, tsBlockBuilder and previousTsBlock + // we have consumed up data from children Operator, just return all remaining cached data in + // cachedTsBlock, tsBlockBuilder and previousTsBlock if (currentIndex >= inputOperatorsCount) { while (previousTsBlock != null) { if (canUseDataFromCachedTsBlock(previousTsBlock)) { @@ -124,14 +130,15 @@ public class LastQuerySortOperator implements ProcessOperator { return res; } - // start stopwatch long maxRuntime = operatorContext.getMaxRunTime().roundTo(TimeUnit.NANOSECONDS); long start = System.nanoTime(); int endIndex = getEndIndex(); - while ((System.nanoTime() - start < maxRuntime) && (currentIndex < endIndex || previousTsBlock != null) && !tsBlockBuilder.isFull()) { + while ((System.nanoTime() - start < maxRuntime) + && (currentIndex < endIndex || previousTsBlock != null) + && !tsBlockBuilder.isFull()) { if (previousTsBlock != null) { if (canUseDataFromCachedTsBlock(previousTsBlock)) { LastQueryUtil.appendLastValue(tsBlockBuilder, cachedTsBlock, cachedTsBlockRowIndex++); @@ -164,7 +171,10 @@ public class LastQuerySortOperator implements ProcessOperator { @Override public boolean hasNext() { - return currentIndex < inputOperatorsCount || cachedTsBlockRowIndex < cachedTsBlockSize || !tsBlockBuilder.isEmpty() || previousTsBlock != null; + return currentIndex < inputOperatorsCount + || cachedTsBlockRowIndex < cachedTsBlockSize + || !tsBlockBuilder.isEmpty() + || previousTsBlock != null; } @Override @@ -185,6 +195,8 @@ public class LastQuerySortOperator implements ProcessOperator { } private boolean canUseDataFromCachedTsBlock(TsBlock tsBlock) { - return cachedTsBlockRowIndex < cachedTsBlockSize && compareTimeSeries(cachedTsBlock, cachedTsBlockRowIndex, tsBlock, 0, timeSeriesComparator) < 0; + return cachedTsBlockRowIndex < cachedTsBlockSize + && compareTimeSeries(cachedTsBlock, cachedTsBlockRowIndex, tsBlock, 0, timeSeriesComparator) + < 0; } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryUtil.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryUtil.java similarity index 86% rename from server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryUtil.java rename to server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryUtil.java index 408e5cd56c..3f08d3a1b2 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryUtil.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/LastQueryUtil.java @@ -16,7 +16,7 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.iotdb.db.mpp.execution.operator; +package org.apache.iotdb.db.mpp.execution.operator.process.last; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.mpp.aggregation.Aggregator; @@ -55,6 +55,10 @@ public class LastQueryUtil { ImmutableList.of(TSDataType.TEXT, TSDataType.TEXT, TSDataType.TEXT)); } + public static Binary getTimeSeries(TsBlock tsBlock, int index) { + return tsBlock.getColumn(0).getBinary(index); + } + public static void appendLastValue( TsBlockBuilder builder, long lastTime, String fullPath, String lastValue, String dataType) { // Time @@ -68,6 +72,18 @@ public class LastQueryUtil { builder.declarePosition(); } + public static void appendLastValue( + TsBlockBuilder builder, long lastTime, Binary fullPath, String lastValue, String dataType) { + // Time + builder.getTimeColumnBuilder().writeLong(lastTime); + // timeseries + builder.getColumnBuilder(0).writeBinary(fullPath); + // value + builder.getColumnBuilder(1).writeBinary(new Binary(lastValue)); + // dataType + builder.getColumnBuilder(2).writeBinary(new Binary(dataType)); + builder.declarePosition(); + } public static void appendLastValue(TsBlockBuilder builder, TsBlock tsBlock) { if (tsBlock.isEmpty()) { @@ -99,7 +115,8 @@ public class LastQueryUtil { builder.declarePosition(); } - public static int compareTimeSeries(TsBlock a, int indexA, TsBlock b, int indexB, Comparator<Binary> comparator) { + public static int compareTimeSeries( + TsBlock a, int indexA, TsBlock b, int indexB, Comparator<Binary> comparator) { return comparator.compare(a.getColumn(0).getBinary(indexA), b.getColumn(0).getBinary(indexB)); } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/UpdateLastCacheOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/UpdateLastCacheOperator.java index 036295ace0..93a5d81cc3 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/UpdateLastCacheOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/last/UpdateLastCacheOperator.java @@ -20,7 +20,6 @@ package org.apache.iotdb.db.mpp.execution.operator.process.last; import org.apache.iotdb.db.metadata.cache.DataNodeSchemaCache; import org.apache.iotdb.db.metadata.path.MeasurementPath; -import org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.Operator; import org.apache.iotdb.db.mpp.execution.operator.OperatorContext; import org.apache.iotdb.db.mpp.execution.operator.process.ProcessOperator; diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java index b8aef10211..cfd8d5f59a 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java @@ -39,7 +39,6 @@ import org.apache.iotdb.db.mpp.execution.exchange.ISourceHandle; import org.apache.iotdb.db.mpp.execution.exchange.MPPDataExchangeManager; import org.apache.iotdb.db.mpp.execution.exchange.MPPDataExchangeService; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext; -import org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.Operator; import org.apache.iotdb.db.mpp.execution.operator.OperatorContext; import org.apache.iotdb.db.mpp.execution.operator.process.AggregationOperator; @@ -47,7 +46,6 @@ import org.apache.iotdb.db.mpp.execution.operator.process.DeviceMergeOperator; import org.apache.iotdb.db.mpp.execution.operator.process.DeviceViewOperator; import org.apache.iotdb.db.mpp.execution.operator.process.FillOperator; import org.apache.iotdb.db.mpp.execution.operator.process.FilterOperator; -import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryOperator; import org.apache.iotdb.db.mpp.execution.operator.process.LimitOperator; import org.apache.iotdb.db.mpp.execution.operator.process.LinearFillOperator; import org.apache.iotdb.db.mpp.execution.operator.process.OffsetOperator; @@ -56,7 +54,6 @@ import org.apache.iotdb.db.mpp.execution.operator.process.RawDataAggregationOper import org.apache.iotdb.db.mpp.execution.operator.process.SlidingWindowAggregationOperator; import org.apache.iotdb.db.mpp.execution.operator.process.TimeJoinOperator; import org.apache.iotdb.db.mpp.execution.operator.process.TransformOperator; -import org.apache.iotdb.db.mpp.execution.operator.process.last.UpdateLastCacheOperator; import org.apache.iotdb.db.mpp.execution.operator.process.fill.IFill; import org.apache.iotdb.db.mpp.execution.operator.process.fill.ILinearFill; import org.apache.iotdb.db.mpp.execution.operator.process.fill.constant.BinaryConstantFill; @@ -77,6 +74,12 @@ import org.apache.iotdb.db.mpp.execution.operator.process.fill.previous.DoublePr import org.apache.iotdb.db.mpp.execution.operator.process.fill.previous.FloatPreviousFill; import org.apache.iotdb.db.mpp.execution.operator.process.fill.previous.IntPreviousFill; import org.apache.iotdb.db.mpp.execution.operator.process.fill.previous.LongPreviousFill; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryCollectOperator; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryMergeOperator; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryOperator; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQuerySortOperator; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil; +import org.apache.iotdb.db.mpp.execution.operator.process.last.UpdateLastCacheOperator; import org.apache.iotdb.db.mpp.execution.operator.process.merge.AscTimeComparator; import org.apache.iotdb.db.mpp.execution.operator.process.merge.ColumnMerger; import org.apache.iotdb.db.mpp.execution.operator.process.merge.DescTimeComparator; @@ -102,7 +105,6 @@ import org.apache.iotdb.db.mpp.execution.operator.source.AlignedSeriesAggregatio import org.apache.iotdb.db.mpp.execution.operator.source.AlignedSeriesScanOperator; import org.apache.iotdb.db.mpp.execution.operator.source.DataSourceOperator; import org.apache.iotdb.db.mpp.execution.operator.source.ExchangeOperator; -import org.apache.iotdb.db.mpp.execution.operator.source.LastCacheScanOperator; import org.apache.iotdb.db.mpp.execution.operator.source.SeriesAggregationScanOperator; import org.apache.iotdb.db.mpp.execution.operator.source.SeriesScanOperator; import org.apache.iotdb.db.mpp.execution.timer.ITimeSliceAllocator; @@ -110,7 +112,6 @@ import org.apache.iotdb.db.mpp.execution.timer.RuleBasedTimeSliceAllocator; import org.apache.iotdb.db.mpp.plan.analyze.TypeProvider; import org.apache.iotdb.db.mpp.plan.expression.leaf.TimeSeriesOperand; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor; import org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.CountSchemaMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.read.DevicesCountNode; @@ -134,13 +135,15 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.ExchangeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FillNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FilterNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.GroupByLevelNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LastQueryMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LimitNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.OffsetNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SlidingWindowAggregationNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SortNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TimeJoinNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TransformNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryCollectNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryMergeNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.sink.FragmentSinkNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedLastQueryScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesAggregationScanNode; @@ -165,12 +168,15 @@ import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder; import org.apache.iotdb.tsfile.read.filter.basic.Filter; import org.apache.iotdb.tsfile.read.filter.operator.Gt; import org.apache.iotdb.tsfile.read.filter.operator.GtEq; +import org.apache.iotdb.tsfile.utils.Binary; +import org.apache.iotdb.tsfile.utils.Pair; import org.apache.commons.lang3.Validate; import java.io.IOException; import java.util.ArrayList; import java.util.Collections; +import java.util.Comparator; import java.util.HashMap; import java.util.HashSet; import java.util.LinkedHashMap; @@ -182,7 +188,7 @@ import java.util.stream.Collectors; import static com.google.common.base.Preconditions.checkArgument; import static java.util.Objects.requireNonNull; -import static org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil.satisfyFilter; +import static org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil.satisfyFilter; import static org.apache.iotdb.db.mpp.plan.constant.DataNodeEndPoints.isSameNode; /** @@ -201,6 +207,10 @@ public class LocalExecutionPlanner { private static final TimeComparator DESC_TIME_COMPARATOR = new DescTimeComparator(); + private static final Comparator<Binary> ASC_BINARY_COMPARATOR = Comparator.naturalOrder(); + + private static final Comparator<Binary> DESC_BINARY_COMPARATOR = Comparator.reverseOrder(); + private static final IdentityFill IDENTITY_FILL = new IdentityFill(); private static final IdentityLinearFill IDENTITY_LINEAR_FILL = new IdentityLinearFill(); @@ -1199,7 +1209,7 @@ public class LocalExecutionPlanner { return null; } } else { // cached last value is satisfied, put it into LastCacheScanOperator - context.addCachedLastValue(timeValuePair, node.getPlanNodeId(), seriesPath.getFullPath()); + context.addCachedLastValue(timeValuePair, seriesPath.getFullPath()); return null; } } @@ -1273,7 +1283,7 @@ public class LocalExecutionPlanner { return null; } } else { // cached last value is satisfied, put it into LastCacheScanOperator - context.addCachedLastValue(timeValuePair, node.getPlanNodeId(), seriesPath.getFullPath()); + context.addCachedLastValue(timeValuePair, seriesPath.getFullPath()); return null; } } @@ -1328,29 +1338,66 @@ public class LocalExecutionPlanner { } @Override - public Operator visitLastQueryMerge( - LastQueryMergeNode node, LocalExecutionPlanContext context) { + public Operator visitLastQuery(LastQueryNode node, LocalExecutionPlanContext context) { + + List<SortItem> sortItemList = node.getMergeOrderParameter().getSortItemList(); + checkArgument( + sortItemList.isEmpty() + || (sortItemList.size() == 1 + && sortItemList.get(0).getSortKey() == SortKey.TIMESERIES), + "Last query only support order by timeseries asc/desc"); context.setLastQueryTimeFilter(node.getTimeFilter()); context.setNeedUpdateLastCache(LastQueryUtil.needUpdateCache(node.getTimeFilter())); - List<Operator> operatorList = + List<UpdateLastCacheOperator> operatorList = node.getChildren().stream() .map(child -> child.accept(this, context)) .filter(Objects::nonNull) + .map(o -> (UpdateLastCacheOperator) o) .collect(Collectors.toList()); - List<TimeValuePair> cachedLastValueList = context.getCachedLastValueList(); + List<Pair<TimeValuePair, Binary>> cachedLastValueAndPathList = + context.getCachedLastValueAndPathList(); - if (cachedLastValueList != null && !cachedLastValueList.isEmpty()) { - TsBlockBuilder builder = LastQueryUtil.createTsBlockBuilder(cachedLastValueList.size()); - for (int i = 0; i < cachedLastValueList.size(); i++) { - TimeValuePair timeValuePair = cachedLastValueList.get(i); - String fullPath = context.cachedLastValuePathList.get(i); + int initSize = cachedLastValueAndPathList != null ? cachedLastValueAndPathList.size() : 0; + // no order by clause + if (sortItemList.isEmpty()) { + TsBlockBuilder builder = LastQueryUtil.createTsBlockBuilder(initSize); + for (int i = 0; i < initSize; i++) { + TimeValuePair timeValuePair = cachedLastValueAndPathList.get(i).left; LastQueryUtil.appendLastValue( builder, timeValuePair.getTimestamp(), - fullPath, + cachedLastValueAndPathList.get(i).right, + timeValuePair.getValue().getStringValue(), + timeValuePair.getValue().getDataType().name()); + } + OperatorContext operatorContext = + context.instanceContext.addOperatorContext( + context.getNextOperatorId(), + node.getPlanNodeId(), + LastQueryOperator.class.getSimpleName()); + context.getTimeSliceAllocator().recordExecutionWeight(operatorContext, 1); + return new LastQueryOperator(operatorContext, operatorList, builder); + } else { + // order by timeseries + Comparator<Binary> comparator = + sortItemList.get(0).getOrdering() == Ordering.ASC + ? ASC_BINARY_COMPARATOR + : DESC_BINARY_COMPARATOR; + // sort values from last cache + if (initSize > 0) { + cachedLastValueAndPathList.sort(Comparator.comparing(Pair::getRight, comparator)); + } + + TsBlockBuilder builder = LastQueryUtil.createTsBlockBuilder(initSize); + for (int i = 0; i < initSize; i++) { + TimeValuePair timeValuePair = cachedLastValueAndPathList.get(i).left; + LastQueryUtil.appendLastValue( + builder, + timeValuePair.getTimestamp(), + cachedLastValueAndPathList.get(i).right, timeValuePair.getValue().getStringValue(), timeValuePair.getValue().getDataType().name()); } @@ -1358,22 +1405,50 @@ public class LocalExecutionPlanner { OperatorContext operatorContext = context.instanceContext.addOperatorContext( context.getNextOperatorId(), - context.firstCachedPlanNodeId, - LastCacheScanOperator.class.getSimpleName()); + node.getPlanNodeId(), + LastQuerySortOperator.class.getSimpleName()); context.getTimeSliceAllocator().recordExecutionWeight(operatorContext, 1); - LastCacheScanOperator operator = - new LastCacheScanOperator( - operatorContext, context.firstCachedPlanNodeId, builder.build()); - operatorList.add(operator); + return new LastQuerySortOperator( + operatorContext, builder.build(), operatorList, comparator); } + } + @Override + public Operator visitLastQueryMerge( + LastQueryMergeNode node, LocalExecutionPlanContext context) { + List<Operator> children = + node.getChildren().stream() + .map(child -> child.accept(this, context)) + .collect(Collectors.toList()); OperatorContext operatorContext = context.instanceContext.addOperatorContext( context.getNextOperatorId(), node.getPlanNodeId(), - LastQueryOperator.class.getSimpleName()); + TimeJoinOperator.class.getSimpleName()); + + SortItem item = node.getMergeOrderParameter().getSortItemList().get(0); + Comparator<Binary> comparator = + item.getOrdering() == Ordering.ASC ? ASC_BINARY_COMPARATOR : DESC_BINARY_COMPARATOR; + + context.getTimeSliceAllocator().recordExecutionWeight(operatorContext, 1); + return new LastQueryMergeOperator(operatorContext, children, comparator); + } + + @Override + public Operator visitLastQueryCollect( + LastQueryCollectNode node, LocalExecutionPlanContext context) { + List<Operator> children = + node.getChildren().stream() + .map(child -> child.accept(this, context)) + .collect(Collectors.toList()); + OperatorContext operatorContext = + context.instanceContext.addOperatorContext( + context.getNextOperatorId(), + node.getPlanNodeId(), + TimeJoinOperator.class.getSimpleName()); + context.getTimeSliceAllocator().recordExecutionWeight(operatorContext, 1); - return new LastQueryOperator(operatorContext, operatorList); + return new LastQueryCollectOperator(operatorContext, children); } private Map<String, List<InputLocation>> makeLayout(PlanNode node) { @@ -1436,13 +1511,9 @@ public class LocalExecutionPlanner { private TypeProvider typeProvider; - // cached last value in last query - private List<TimeValuePair> cachedLastValueList; - // full path for each cached last value, this size should be equal to cachedLastValueList - private List<String> cachedLastValuePathList; - // PlanNodeId of first LastQueryScanNode/AlignedLastQueryScanNode, it's used for sourceId of - // LastCachedScanOperator - private PlanNodeId firstCachedPlanNodeId; + // left is cached last value in last query + // right is full path for each cached last value + private List<Pair<TimeValuePair, Binary>> cachedLastValueAndPathList; // timeFilter for last query private Filter lastQueryTimeFilter; // whether we need to update last cache @@ -1504,19 +1575,15 @@ public class LocalExecutionPlanner { this.needUpdateLastCache = needUpdateLastCache; } - public void addCachedLastValue( - TimeValuePair timeValuePair, PlanNodeId planNodeId, String fullPath) { - if (cachedLastValueList == null) { - cachedLastValueList = new ArrayList<>(); - cachedLastValuePathList = new ArrayList<>(); - firstCachedPlanNodeId = planNodeId; + public void addCachedLastValue(TimeValuePair timeValuePair, String fullPath) { + if (cachedLastValueAndPathList == null) { + cachedLastValueAndPathList = new ArrayList<>(); } - cachedLastValueList.add(timeValuePair); - cachedLastValuePathList.add(fullPath); + cachedLastValueAndPathList.add(new Pair<>(timeValuePair, new Binary(fullPath))); } - public List<TimeValuePair> getCachedLastValueList() { - return cachedLastValueList; + public List<Pair<TimeValuePair, Binary>> getCachedLastValueAndPathList() { + return cachedLastValueAndPathList; } public ISinkHandle getSinkHandle() { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java index c1dd00a1f2..016cfeb381 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanBuilder.java @@ -52,12 +52,12 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceViewNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FillNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FilterNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.GroupByLevelNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LastQueryMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LimitNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.OffsetNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SlidingWindowAggregationNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TimeJoinNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TransformNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedLastQueryScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesAggregationScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesScanNode; @@ -160,7 +160,7 @@ public class LogicalPlanBuilder { } this.root = - new LastQueryMergeNode( + new LastQueryNode( context.getQueryId().genPlanNodeId(), sourceNodeList, globalTimeFilter, diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java index 4159db3a7e..b375cf3739 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/ExchangeNodeAdder.java @@ -36,11 +36,11 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceViewNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.ExchangeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.GroupByLevelNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LastQueryMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.MultiChildNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SlidingWindowAggregationNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TimeJoinNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TransformNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedLastQueryScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesAggregationScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesScanNode; @@ -188,7 +188,7 @@ public class ExchangeNodeAdder extends PlanVisitor<PlanNode, NodeGroupContext> { } @Override - public PlanNode visitLastQueryMerge(LastQueryMergeNode node, NodeGroupContext context) { + public PlanNode visitLastQuery(LastQueryNode node, NodeGroupContext context) { return processMultiChildNode(node, context); } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java index 51618cbe64..179d5a600c 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SourceRewriter.java @@ -36,10 +36,10 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.AggregationNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceViewNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.GroupByLevelNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LastQueryMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.MultiChildNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SlidingWindowAggregationNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TimeJoinNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedLastQueryScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesAggregationScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesScanNode; @@ -230,8 +230,8 @@ public class SourceRewriter extends SimplePlanNodeRewriter<DistributionPlanConte @Override public PlanNode visitLastQueryScan(LastQueryScanNode node, DistributionPlanContext context) { - LastQueryMergeNode mergeNode = - new LastQueryMergeNode( + LastQueryNode mergeNode = + new LastQueryNode( context.queryContext.getQueryId().genPlanNodeId(), node.getPartitionTimeFilter(), new OrderByParameter()); @@ -241,8 +241,8 @@ public class SourceRewriter extends SimplePlanNodeRewriter<DistributionPlanConte @Override public PlanNode visitAlignedLastQueryScan( AlignedLastQueryScanNode node, DistributionPlanContext context) { - LastQueryMergeNode mergeNode = - new LastQueryMergeNode( + LastQueryNode mergeNode = + new LastQueryNode( context.queryContext.getQueryId().genPlanNodeId(), node.getPartitionTimeFilter(), new OrderByParameter()); @@ -370,7 +370,7 @@ public class SourceRewriter extends SimplePlanNodeRewriter<DistributionPlanConte } @Override - public PlanNode visitLastQueryMerge(LastQueryMergeNode node, DistributionPlanContext context) { + public PlanNode visitLastQuery(LastQueryNode node, DistributionPlanContext context) { // For last query, we need to keep every FI's root node is LastQueryMergeNode. So we // force every region group have a parent node even if there is only 1 child for it. context.setForceAddParent(true); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java index 4c176d83b2..c20b994ea8 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java @@ -46,7 +46,6 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.ExchangeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FillNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FilterNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.GroupByLevelNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LastQueryMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LimitNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.OffsetNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.ProjectNode; @@ -54,6 +53,9 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SlidingWindowAggre import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SortNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TimeJoinNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TransformNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryCollectNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryMergeNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.sink.FragmentSinkNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedLastQueryScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesAggregationScanNode; @@ -122,9 +124,11 @@ public enum PlanNodeType { DELETE_TIMESERIES((short) 45), LAST_QUERY_SCAN((short) 46), ALIGNED_LAST_QUERY_SCAN((short) 47), - LAST_QUERY_MERGE((short) 48), - NODE_PATHS_COUNT((short) 49), - INTERNAL_CREATE_TIMESERIES((short) 50); + LAST_QUERY((short) 48), + LAST_QUERY_MERGE((short) 49), + LAST_QUERY_COLLECT((short) 50), + NODE_PATHS_COUNT((short) 51), + INTERNAL_CREATE_TIMESERIES((short) 52); public static final int BYTES = Short.BYTES; @@ -262,10 +266,14 @@ public enum PlanNodeType { case 47: return AlignedLastQueryScanNode.deserialize(buffer); case 48: - return LastQueryMergeNode.deserialize(buffer); + return LastQueryNode.deserialize(buffer); case 49: - return NodePathsCountNode.deserialize(buffer); + return LastQueryMergeNode.deserialize(buffer); case 50: + return LastQueryCollectNode.deserialize(buffer); + case 51: + return NodePathsCountNode.deserialize(buffer); + case 52: return InternalCreateTimeSeriesNode.deserialize(buffer); default: throw new IllegalArgumentException("Invalid node type: " + nodeType); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java index fc83e2a5a7..907e7985e6 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java @@ -46,7 +46,6 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.ExchangeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FillNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FilterNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.GroupByLevelNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LastQueryMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LimitNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.OffsetNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.ProjectNode; @@ -54,6 +53,9 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SlidingWindowAggre import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.SortNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TimeJoinNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TransformNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryCollectNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryMergeNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.sink.FragmentSinkNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedLastQueryScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesAggregationScanNode; @@ -256,10 +258,18 @@ public abstract class PlanVisitor<R, C> { return visitPlan(node, context); } + public R visitLastQuery(LastQueryNode node, C context) { + return visitPlan(node, context); + } + public R visitLastQueryMerge(LastQueryMergeNode node, C context) { return visitPlan(node, context); } + public R visitLastQueryCollect(LastQueryCollectNode node, C context) { + return visitPlan(node, context); + } + public R visitDeleteTimeseries(DeleteTimeSeriesNode node, C context) { return visitPlan(node, context); } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryCollectNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryCollectNode.java new file mode 100644 index 0000000000..a81a5fb455 --- /dev/null +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryCollectNode.java @@ -0,0 +1,108 @@ +/* + * 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.mpp.plan.planner.plan.node.process.last; + +import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeType; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.MultiChildNode; + +import java.io.DataOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.util.List; +import java.util.Objects; + +import static org.apache.iotdb.db.mpp.plan.planner.plan.node.source.LastQueryScanNode.LAST_QUERY_HEADER_COLUMNS; + +public class LastQueryCollectNode extends MultiChildNode { + + public LastQueryCollectNode(PlanNodeId id) { + super(id); + } + + public LastQueryCollectNode(PlanNodeId id, List<PlanNode> children) { + super(id, children); + } + + @Override + public List<PlanNode> getChildren() { + return children; + } + + @Override + public void addChild(PlanNode child) { + children.add(child); + } + + @Override + public PlanNode clone() { + return new LastQueryCollectNode(getPlanNodeId()); + } + + @Override + public int allowedChildCount() { + return CHILD_COUNT_NO_LIMIT; + } + + @Override + public List<String> getOutputColumnNames() { + return LAST_QUERY_HEADER_COLUMNS; + } + + @Override + public String toString() { + return String.format("LastQueryCollectNode-%s", this.getPlanNodeId()); + } + + @Override + public boolean equals(Object o) { + return super.equals(o); + } + + @Override + public int hashCode() { + return Objects.hash(super.hashCode()); + } + + @Override + public <R, C> R accept(PlanVisitor<R, C> visitor, C context) { + return visitor.visitLastQueryCollect(this, context); + } + + @Override + protected void serializeAttributes(ByteBuffer byteBuffer) { + PlanNodeType.LAST_QUERY_COLLECT.serialize(byteBuffer); + } + + @Override + protected void serializeAttributes(DataOutputStream stream) throws IOException { + PlanNodeType.LAST_QUERY_COLLECT.serialize(stream); + } + + public static LastQueryCollectNode deserialize(ByteBuffer byteBuffer) { + PlanNodeId planNodeId = PlanNodeId.deserialize(byteBuffer); + return new LastQueryCollectNode(planNodeId); + } + + public void setChildren(List<PlanNode> children) { + this.children = children; + } +} diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/LastQueryMergeNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryMergeNode.java similarity index 68% copy from server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/LastQueryMergeNode.java copy to server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryMergeNode.java index fd1dd2afba..482fc0110d 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/LastQueryMergeNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryMergeNode.java @@ -16,18 +16,14 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.iotdb.db.mpp.plan.planner.plan.node.process; +package org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeType; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.MultiChildNode; import org.apache.iotdb.db.mpp.plan.planner.plan.parameter.OrderByParameter; -import org.apache.iotdb.tsfile.read.filter.basic.Filter; -import org.apache.iotdb.tsfile.read.filter.factory.FilterFactory; -import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils; - -import javax.annotation.Nullable; import java.io.DataOutputStream; import java.io.IOException; @@ -39,26 +35,18 @@ import static org.apache.iotdb.db.mpp.plan.planner.plan.node.source.LastQuerySca public class LastQueryMergeNode extends MultiChildNode { - private final Filter timeFilter; - // The result output order, which could sort by sensor and time. // The size of this list is 2 and the first SortItem in this list has higher priority. private final OrderByParameter mergeOrderParameter; - public LastQueryMergeNode( - PlanNodeId id, Filter timeFilter, OrderByParameter mergeOrderParameter) { + public LastQueryMergeNode(PlanNodeId id, OrderByParameter mergeOrderParameter) { super(id); - this.timeFilter = timeFilter; this.mergeOrderParameter = mergeOrderParameter; } public LastQueryMergeNode( - PlanNodeId id, - List<PlanNode> children, - Filter timeFilter, - OrderByParameter mergeOrderParameter) { + PlanNodeId id, List<PlanNode> children, OrderByParameter mergeOrderParameter) { super(id, children); - this.timeFilter = timeFilter; this.mergeOrderParameter = mergeOrderParameter; } @@ -74,7 +62,7 @@ public class LastQueryMergeNode extends MultiChildNode { @Override public PlanNode clone() { - return new LastQueryMergeNode(getPlanNodeId(), timeFilter, mergeOrderParameter); + return new LastQueryMergeNode(getPlanNodeId(), mergeOrderParameter); } @Override @@ -90,7 +78,7 @@ public class LastQueryMergeNode extends MultiChildNode { @Override public String toString() { return String.format( - "LastQueryMergeNode-%s:[TimeFilter: %s]", this.getPlanNodeId(), timeFilter); + "LastQueryMergeNode-%s:[OrderByParameter: %s]", this.getPlanNodeId(), mergeOrderParameter); } @Override @@ -105,13 +93,12 @@ public class LastQueryMergeNode extends MultiChildNode { return false; } LastQueryMergeNode that = (LastQueryMergeNode) o; - return Objects.equals(timeFilter, that.timeFilter) - && mergeOrderParameter.equals(that.mergeOrderParameter); + return mergeOrderParameter.equals(that.mergeOrderParameter); } @Override public int hashCode() { - return Objects.hash(super.hashCode(), timeFilter, mergeOrderParameter); + return Objects.hash(super.hashCode(), mergeOrderParameter); } @Override @@ -122,43 +109,26 @@ public class LastQueryMergeNode extends MultiChildNode { @Override protected void serializeAttributes(ByteBuffer byteBuffer) { PlanNodeType.LAST_QUERY_MERGE.serialize(byteBuffer); - if (timeFilter == null) { - ReadWriteIOUtils.write((byte) 0, byteBuffer); - } else { - ReadWriteIOUtils.write((byte) 1, byteBuffer); - timeFilter.serialize(byteBuffer); - } mergeOrderParameter.serializeAttributes(byteBuffer); } @Override protected void serializeAttributes(DataOutputStream stream) throws IOException { PlanNodeType.LAST_QUERY_MERGE.serialize(stream); - if (timeFilter == null) { - ReadWriteIOUtils.write((byte) 0, stream); - } else { - ReadWriteIOUtils.write((byte) 1, stream); - timeFilter.serialize(stream); - } mergeOrderParameter.serializeAttributes(stream); } public static LastQueryMergeNode deserialize(ByteBuffer byteBuffer) { - Filter timeFilter = null; - if (ReadWriteIOUtils.readByte(byteBuffer) == 1) { - timeFilter = FilterFactory.deserialize(byteBuffer); - } OrderByParameter mergeOrderParameter = OrderByParameter.deserialize(byteBuffer); PlanNodeId planNodeId = PlanNodeId.deserialize(byteBuffer); - return new LastQueryMergeNode(planNodeId, timeFilter, mergeOrderParameter); + return new LastQueryMergeNode(planNodeId, mergeOrderParameter); } public void setChildren(List<PlanNode> children) { this.children = children; } - @Nullable - public Filter getTimeFilter() { - return timeFilter; + public OrderByParameter getMergeOrderParameter() { + return mergeOrderParameter; } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/LastQueryMergeNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryNode.java similarity index 84% rename from server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/LastQueryMergeNode.java rename to server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryNode.java index fd1dd2afba..6eb7256b91 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/LastQueryMergeNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/process/last/LastQueryNode.java @@ -16,12 +16,13 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.iotdb.db.mpp.plan.planner.plan.node.process; +package org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeType; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.MultiChildNode; import org.apache.iotdb.db.mpp.plan.planner.plan.parameter.OrderByParameter; import org.apache.iotdb.tsfile.read.filter.basic.Filter; import org.apache.iotdb.tsfile.read.filter.factory.FilterFactory; @@ -37,7 +38,7 @@ import java.util.Objects; import static org.apache.iotdb.db.mpp.plan.planner.plan.node.source.LastQueryScanNode.LAST_QUERY_HEADER_COLUMNS; -public class LastQueryMergeNode extends MultiChildNode { +public class LastQueryNode extends MultiChildNode { private final Filter timeFilter; @@ -45,14 +46,13 @@ public class LastQueryMergeNode extends MultiChildNode { // The size of this list is 2 and the first SortItem in this list has higher priority. private final OrderByParameter mergeOrderParameter; - public LastQueryMergeNode( - PlanNodeId id, Filter timeFilter, OrderByParameter mergeOrderParameter) { + public LastQueryNode(PlanNodeId id, Filter timeFilter, OrderByParameter mergeOrderParameter) { super(id); this.timeFilter = timeFilter; this.mergeOrderParameter = mergeOrderParameter; } - public LastQueryMergeNode( + public LastQueryNode( PlanNodeId id, List<PlanNode> children, Filter timeFilter, @@ -74,7 +74,7 @@ public class LastQueryMergeNode extends MultiChildNode { @Override public PlanNode clone() { - return new LastQueryMergeNode(getPlanNodeId(), timeFilter, mergeOrderParameter); + return new LastQueryNode(getPlanNodeId(), timeFilter, mergeOrderParameter); } @Override @@ -104,7 +104,7 @@ public class LastQueryMergeNode extends MultiChildNode { if (!super.equals(o)) { return false; } - LastQueryMergeNode that = (LastQueryMergeNode) o; + LastQueryNode that = (LastQueryNode) o; return Objects.equals(timeFilter, that.timeFilter) && mergeOrderParameter.equals(that.mergeOrderParameter); } @@ -116,12 +116,12 @@ public class LastQueryMergeNode extends MultiChildNode { @Override public <R, C> R accept(PlanVisitor<R, C> visitor, C context) { - return visitor.visitLastQueryMerge(this, context); + return visitor.visitLastQuery(this, context); } @Override protected void serializeAttributes(ByteBuffer byteBuffer) { - PlanNodeType.LAST_QUERY_MERGE.serialize(byteBuffer); + PlanNodeType.LAST_QUERY.serialize(byteBuffer); if (timeFilter == null) { ReadWriteIOUtils.write((byte) 0, byteBuffer); } else { @@ -133,7 +133,7 @@ public class LastQueryMergeNode extends MultiChildNode { @Override protected void serializeAttributes(DataOutputStream stream) throws IOException { - PlanNodeType.LAST_QUERY_MERGE.serialize(stream); + PlanNodeType.LAST_QUERY.serialize(stream); if (timeFilter == null) { ReadWriteIOUtils.write((byte) 0, stream); } else { @@ -143,14 +143,14 @@ public class LastQueryMergeNode extends MultiChildNode { mergeOrderParameter.serializeAttributes(stream); } - public static LastQueryMergeNode deserialize(ByteBuffer byteBuffer) { + public static LastQueryNode deserialize(ByteBuffer byteBuffer) { Filter timeFilter = null; if (ReadWriteIOUtils.readByte(byteBuffer) == 1) { timeFilter = FilterFactory.deserialize(byteBuffer); } OrderByParameter mergeOrderParameter = OrderByParameter.deserialize(byteBuffer); PlanNodeId planNodeId = PlanNodeId.deserialize(byteBuffer); - return new LastQueryMergeNode(planNodeId, timeFilter, mergeOrderParameter); + return new LastQueryNode(planNodeId, timeFilter, mergeOrderParameter); } public void setChildren(List<PlanNode> children) { @@ -161,4 +161,8 @@ public class LastQueryMergeNode extends MultiChildNode { public Filter getTimeFilter() { return timeFilter; } + + public OrderByParameter getMergeOrderParameter() { + return mergeOrderParameter; + } } diff --git a/server/src/main/java/org/apache/iotdb/db/query/executor/LastQueryExecutor.java b/server/src/main/java/org/apache/iotdb/db/query/executor/LastQueryExecutor.java index af9526b142..76ce6ff155 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/executor/LastQueryExecutor.java +++ b/server/src/main/java/org/apache/iotdb/db/query/executor/LastQueryExecutor.java @@ -67,7 +67,7 @@ import java.util.stream.Collectors; import static org.apache.iotdb.commons.conf.IoTDBConstant.COLUMN_TIMESERIES; import static org.apache.iotdb.commons.conf.IoTDBConstant.COLUMN_TIMESERIES_DATATYPE; import static org.apache.iotdb.commons.conf.IoTDBConstant.COLUMN_VALUE; -import static org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil.satisfyFilter; +import static org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil.satisfyFilter; public class LastQueryExecutor { diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/AggregationOperatorTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/AggregationOperatorTest.java index 69b470845b..233965aa3e 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/AggregationOperatorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/AggregationOperatorTest.java @@ -65,7 +65,7 @@ import static org.junit.Assert.assertEquals; public class AggregationOperatorTest { - public static Duration TEST_TIME_SLICE = new Duration(500, TimeUnit.MILLISECONDS); + public static Duration TEST_TIME_SLICE = new Duration(50000, TimeUnit.MILLISECONDS); private static final String AGGREGATION_OPERATOR_TEST_SG = "root.AggregationOperatorTest"; private final List<String> deviceIds = new ArrayList<>(); diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastCacheScanOperatorTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastCacheScanOperatorTest.java deleted file mode 100644 index 72ee86fb65..0000000000 --- a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastCacheScanOperatorTest.java +++ /dev/null @@ -1,93 +0,0 @@ -/* - * 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.mpp.execution.operator; - -import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; -import org.apache.iotdb.db.mpp.common.FragmentInstanceId; -import org.apache.iotdb.db.mpp.common.PlanFragmentId; -import org.apache.iotdb.db.mpp.common.QueryId; -import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext; -import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceStateMachine; -import org.apache.iotdb.db.mpp.execution.operator.source.LastCacheScanOperator; -import org.apache.iotdb.db.mpp.execution.operator.source.SeriesScanOperator; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; -import org.apache.iotdb.tsfile.read.common.block.TsBlock; -import org.apache.iotdb.tsfile.read.common.block.TsBlockBuilder; - -import org.junit.Test; - -import java.util.concurrent.ExecutorService; - -import static org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext.createFragmentInstanceContext; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertTrue; - -public class LastCacheScanOperatorTest { - - @Test - public void batchTest() { - ExecutorService instanceNotificationExecutor = - IoTDBThreadPoolFactory.newFixedThreadPool(1, "test-instance-notification"); - try { - QueryId queryId = new QueryId("stub_query"); - FragmentInstanceId instanceId = - new FragmentInstanceId(new PlanFragmentId(queryId, 0), "stub-instance"); - FragmentInstanceStateMachine stateMachine = - new FragmentInstanceStateMachine(instanceId, instanceNotificationExecutor); - FragmentInstanceContext fragmentInstanceContext = - createFragmentInstanceContext(instanceId, stateMachine); - PlanNodeId planNodeId1 = new PlanNodeId("1"); - fragmentInstanceContext.addOperatorContext( - 1, planNodeId1, SeriesScanOperator.class.getSimpleName()); - - TsBlockBuilder builder = LastQueryUtil.createTsBlockBuilder(6); - - LastQueryUtil.appendLastValue(builder, 1, "root.sg.d.s1", "true", "BOOLEAN"); - LastQueryUtil.appendLastValue(builder, 2, "root.sg.d.s2", "2", "INT32"); - LastQueryUtil.appendLastValue(builder, 3, "root.sg.d.s3", "3", "INT64"); - LastQueryUtil.appendLastValue(builder, 4, "root.sg.d.s4", "4.4", "FLOAT"); - LastQueryUtil.appendLastValue(builder, 3, "root.sg.d.s5", "3.3", "DOUBLE"); - LastQueryUtil.appendLastValue(builder, 1, "root.sg.d.s6", "peace", "TEXT"); - - TsBlock tsBlock = builder.build(); - - LastCacheScanOperator lastCacheScanOperator = - new LastCacheScanOperator( - fragmentInstanceContext.getOperatorContexts().get(0), planNodeId1, tsBlock); - - assertTrue(lastCacheScanOperator.isBlocked().isDone()); - assertTrue(lastCacheScanOperator.hasNext()); - TsBlock result = lastCacheScanOperator.next(); - assertEquals(tsBlock.getPositionCount(), result.getPositionCount()); - assertEquals(tsBlock.getValueColumnCount(), result.getValueColumnCount()); - for (int i = 0; i < tsBlock.getPositionCount(); i++) { - assertEquals(tsBlock.getTimeByIndex(i), result.getTimeByIndex(i)); - for (int j = 0; j < tsBlock.getValueColumnCount(); j++) { - assertEquals(tsBlock.getColumn(j).getBinary(i), result.getColumn(j).getBinary(i)); - } - } - assertFalse(lastCacheScanOperator.hasNext()); - assertTrue(lastCacheScanOperator.isFinished()); - - } finally { - instanceNotificationExecutor.shutdown(); - } - } -} diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryOperatorTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryOperatorTest.java index 5804095ffd..97bfc25e8c 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryOperatorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryOperatorTest.java @@ -31,8 +31,8 @@ import org.apache.iotdb.db.mpp.common.QueryId; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceStateMachine; import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryOperator; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.process.last.UpdateLastCacheOperator; -import org.apache.iotdb.db.mpp.execution.operator.source.LastCacheScanOperator; import org.apache.iotdb.db.mpp.execution.operator.source.SeriesAggregationScanOperator; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.query.reader.series.SeriesReaderTestUtil; @@ -62,7 +62,7 @@ import static org.junit.Assert.fail; public class LastQueryOperatorTest { - private static final String SERIES_SCAN_OPERATOR_TEST_SG = "root.LastQueryMergeOperatorTest"; + private static final String SERIES_SCAN_OPERATOR_TEST_SG = "root.LastQueryOperatorTest"; private final List<String> deviceIds = new ArrayList<>(); private final List<MeasurementSchema> measurementSchemas = new ArrayList<>(); @@ -173,7 +173,8 @@ public class LastQueryOperatorTest { LastQueryOperator lastQueryOperator = new LastQueryOperator( fragmentInstanceContext.getOperatorContexts().get(4), - ImmutableList.of(updateLastCacheOperator1, updateLastCacheOperator2)); + ImmutableList.of(updateLastCacheOperator1, updateLastCacheOperator2), + LastQueryUtil.createTsBlockBuilder()); int count = 0; while (!lastQueryOperator.isFinished()) { @@ -235,11 +236,7 @@ public class LastQueryOperatorTest { PlanNodeId planNodeId5 = new PlanNodeId("5"); fragmentInstanceContext.addOperatorContext( - 5, planNodeId4, LastCacheScanOperator.class.getSimpleName()); - - PlanNodeId planNodeId6 = new PlanNodeId("6"); - fragmentInstanceContext.addOperatorContext( - 6, planNodeId6, LastQueryOperator.class.getSimpleName()); + 5, planNodeId4, LastQueryOperator.class.getSimpleName()); fragmentInstanceContext .getOperatorContexts() @@ -301,19 +298,14 @@ public class LastQueryOperatorTest { LastQueryUtil.appendLastValue( builder, 499, SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor4", "10499", "INT32"); - TsBlock tsBlock = builder.build(); - - LastCacheScanOperator lastCacheScanOperator = - new LastCacheScanOperator( - fragmentInstanceContext.getOperatorContexts().get(4), planNodeId5, tsBlock); - LastQueryOperator lastQueryOperator = new LastQueryOperator( - fragmentInstanceContext.getOperatorContexts().get(5), - ImmutableList.of( - updateLastCacheOperator1, updateLastCacheOperator2, lastCacheScanOperator)); + fragmentInstanceContext.getOperatorContexts().get(4), + ImmutableList.of(updateLastCacheOperator1, updateLastCacheOperator2), + builder); int count = 0; + int[] suffix = new int[] {2, 3, 4, 0, 1}; while (!lastQueryOperator.isFinished()) { assertTrue(lastQueryOperator.isBlocked().isDone()); assertTrue(lastQueryOperator.hasNext()); @@ -326,7 +318,7 @@ public class LastQueryOperatorTest { for (int i = 0; i < result.getPositionCount(); i++) { assertEquals(499, result.getTimeByIndex(i)); assertEquals( - SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor" + count, + SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor" + suffix[count], result.getColumn(0).getBinary(i).toString()); assertEquals("10499", result.getColumn(1).getBinary(i).toString()); assertEquals(TSDataType.INT32.name(), result.getColumn(2).getBinary(i).toString()); diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryOperatorTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQuerySortOperatorTest.java similarity index 90% copy from server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryOperatorTest.java copy to server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQuerySortOperatorTest.java index 5804095ffd..4f575d642d 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQueryOperatorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/LastQuerySortOperatorTest.java @@ -31,8 +31,9 @@ import org.apache.iotdb.db.mpp.common.QueryId; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceStateMachine; import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryOperator; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQuerySortOperator; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.process.last.UpdateLastCacheOperator; -import org.apache.iotdb.db.mpp.execution.operator.source.LastCacheScanOperator; import org.apache.iotdb.db.mpp.execution.operator.source.SeriesAggregationScanOperator; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.query.reader.series.SeriesReaderTestUtil; @@ -50,6 +51,7 @@ import org.junit.Test; import java.io.IOException; import java.util.ArrayList; +import java.util.Comparator; import java.util.List; import java.util.Set; import java.util.concurrent.ExecutorService; @@ -60,9 +62,9 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; -public class LastQueryOperatorTest { +public class LastQuerySortOperatorTest { - private static final String SERIES_SCAN_OPERATOR_TEST_SG = "root.LastQueryMergeOperatorTest"; + private static final String SERIES_SCAN_OPERATOR_TEST_SG = "root.LastQuerySortOperatorTest"; private final List<String> deviceIds = new ArrayList<>(); private final List<MeasurementSchema> measurementSchemas = new ArrayList<>(); @@ -170,16 +172,18 @@ public class LastQueryOperatorTest { null, false); - LastQueryOperator lastQueryOperator = - new LastQueryOperator( + LastQuerySortOperator lastQuerySortOperator = + new LastQuerySortOperator( fragmentInstanceContext.getOperatorContexts().get(4), - ImmutableList.of(updateLastCacheOperator1, updateLastCacheOperator2)); + LastQueryUtil.createTsBlockBuilder().build(), + ImmutableList.of(updateLastCacheOperator1, updateLastCacheOperator2), + Comparator.naturalOrder()); int count = 0; - while (!lastQueryOperator.isFinished()) { - assertTrue(lastQueryOperator.isBlocked().isDone()); - assertTrue(lastQueryOperator.hasNext()); - TsBlock result = lastQueryOperator.next(); + while (!lastQuerySortOperator.isFinished()) { + assertTrue(lastQuerySortOperator.isBlocked().isDone()); + assertTrue(lastQuerySortOperator.hasNext()); + TsBlock result = lastQuerySortOperator.next(); if (result == null) { continue; } @@ -235,11 +239,7 @@ public class LastQueryOperatorTest { PlanNodeId planNodeId5 = new PlanNodeId("5"); fragmentInstanceContext.addOperatorContext( - 5, planNodeId4, LastCacheScanOperator.class.getSimpleName()); - - PlanNodeId planNodeId6 = new PlanNodeId("6"); - fragmentInstanceContext.addOperatorContext( - 6, planNodeId6, LastQueryOperator.class.getSimpleName()); + 5, planNodeId4, LastQueryOperator.class.getSimpleName()); fragmentInstanceContext .getOperatorContexts() @@ -295,29 +295,25 @@ public class LastQueryOperatorTest { TsBlockBuilder builder = LastQueryUtil.createTsBlockBuilder(6); LastQueryUtil.appendLastValue( - builder, 499, SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor2", "10499", "INT32"); + builder, 499, SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor4", "10499", "INT32"); LastQueryUtil.appendLastValue( builder, 499, SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor3", "10499", "INT32"); LastQueryUtil.appendLastValue( - builder, 499, SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor4", "10499", "INT32"); - - TsBlock tsBlock = builder.build(); - - LastCacheScanOperator lastCacheScanOperator = - new LastCacheScanOperator( - fragmentInstanceContext.getOperatorContexts().get(4), planNodeId5, tsBlock); + builder, 499, SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor2", "10499", "INT32"); - LastQueryOperator lastQueryOperator = - new LastQueryOperator( - fragmentInstanceContext.getOperatorContexts().get(5), - ImmutableList.of( - updateLastCacheOperator1, updateLastCacheOperator2, lastCacheScanOperator)); + LastQuerySortOperator lastQuerySortOperator = + new LastQuerySortOperator( + fragmentInstanceContext.getOperatorContexts().get(4), + builder.build(), + ImmutableList.of(updateLastCacheOperator2, updateLastCacheOperator1), + Comparator.reverseOrder()); int count = 0; - while (!lastQueryOperator.isFinished()) { - assertTrue(lastQueryOperator.isBlocked().isDone()); - assertTrue(lastQueryOperator.hasNext()); - TsBlock result = lastQueryOperator.next(); + int[] suffix = new int[] {4, 3, 2, 1, 0}; + while (!lastQuerySortOperator.isFinished()) { + assertTrue(lastQuerySortOperator.isBlocked().isDone()); + assertTrue(lastQuerySortOperator.hasNext()); + TsBlock result = lastQuerySortOperator.next(); if (result == null) { continue; } @@ -326,7 +322,7 @@ public class LastQueryOperatorTest { for (int i = 0; i < result.getPositionCount(); i++) { assertEquals(499, result.getTimeByIndex(i)); assertEquals( - SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor" + count, + SERIES_SCAN_OPERATOR_TEST_SG + ".device0.sensor" + suffix[count], result.getColumn(0).getBinary(i).toString()); assertEquals("10499", result.getColumn(1).getBinary(i).toString()); assertEquals(TSDataType.INT32.name(), result.getColumn(2).getBinary(i).toString()); diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/UpdateLastCacheOperatorTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/UpdateLastCacheOperatorTest.java index 4567680edb..34b9b8fa87 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/UpdateLastCacheOperatorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/UpdateLastCacheOperatorTest.java @@ -30,6 +30,7 @@ import org.apache.iotdb.db.mpp.common.PlanFragmentId; import org.apache.iotdb.db.mpp.common.QueryId; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceStateMachine; +import org.apache.iotdb.db.mpp.execution.operator.process.last.LastQueryUtil; import org.apache.iotdb.db.mpp.execution.operator.process.last.UpdateLastCacheOperator; import org.apache.iotdb.db.mpp.execution.operator.source.SeriesAggregationScanOperator; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/QueryLogicalPlanUtil.java b/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/QueryLogicalPlanUtil.java index 871fda2aed..7e9622ddb4 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/QueryLogicalPlanUtil.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/plan/plan/QueryLogicalPlanUtil.java @@ -36,10 +36,10 @@ import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.AggregationNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceViewNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.FilterNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.GroupByLevelNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LastQueryMergeNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.LimitNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.OffsetNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.TimeJoinNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.last.LastQueryNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedLastQueryScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesAggregationScanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.source.AlignedSeriesScanNode; @@ -140,8 +140,8 @@ public class QueryLogicalPlanUtil { new LastQueryScanNode( queryId.genPlanNodeId(), (MeasurementPath) schemaMap.get("root.sg.d2.s2"))); - LastQueryMergeNode lastQueryMergeNode = - new LastQueryMergeNode( + LastQueryNode lastQueryNode = + new LastQueryNode( queryId.genPlanNodeId(), sourceNodeList, TimeFilter.gt(100), @@ -149,7 +149,7 @@ public class QueryLogicalPlanUtil { Collections.singletonList(new SortItem(SortKey.TIMESERIES, Ordering.ASC)))); querySQLs.add(sql); - sqlToPlanMap.put(sql, lastQueryMergeNode); + sqlToPlanMap.put(sql, lastQueryNode); } /* Simple Query */
