This is an automated email from the ASF dual-hosted git repository. xiangweiwei pushed a commit to branch aggControlLogic in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit d6d28fcfea008b1571d6a12cfaa9002aff4effeb Author: Alima777 <[email protected]> AuthorDate: Tue May 31 17:13:43 2022 +0800 modify SlidingWindowAggregationOperator control next logic --- .../process/RawDataAggregationOperator.java | 4 +-- .../process/SlidingWindowAggregationOperator.java | 29 +++++++++++++++------- .../SlidingWindowAggregationOperatorTest.java | 7 +++++- 3 files changed, 28 insertions(+), 12 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/RawDataAggregationOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/RawDataAggregationOperator.java index 683d32fbac..9581975616 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/RawDataAggregationOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/RawDataAggregationOperator.java @@ -106,13 +106,13 @@ public class RawDataAggregationOperator implements ProcessOperator { // 1. Calculate aggregation result based on current time window boolean canCallNext = true; while (!calcFromCacheData(curTimeRange)) { + preCachedData = null; // child.next can only be invoked once if (child.hasNext() && canCallNext) { - canCallNext = false; preCachedData = child.next(); + canCallNext = false; // if child still has next but can't be invoked now } else if (child.hasNext()) { - preCachedData = null; return null; } else { break; diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/SlidingWindowAggregationOperator.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/SlidingWindowAggregationOperator.java index 150699ccb3..6a8288bd4d 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/SlidingWindowAggregationOperator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/process/SlidingWindowAggregationOperator.java @@ -52,6 +52,8 @@ public class SlidingWindowAggregationOperator implements ProcessOperator { private final List<SlidingWindowAggregator> aggregators; private final ITimeRangeIterator timeRangeIterator; + // current interval of aggregation window [curStartTime, curEndTime) + private TimeRange curTimeRange; private final boolean ascending; @@ -81,28 +83,37 @@ public class SlidingWindowAggregationOperator implements ProcessOperator { @Override public boolean hasNext() { - return timeRangeIterator.hasNextTimeRange(); + return curTimeRange != null || timeRangeIterator.hasNextTimeRange(); } @Override public TsBlock next() { - // 1. Clear previous aggregation result - TimeRange curTimeRange = timeRangeIterator.nextTimeRange(); - for (SlidingWindowAggregator aggregator : aggregators) { - aggregator.updateTimeRange(curTimeRange); + // Move to next timeRange + if (curTimeRange == null && timeRangeIterator.hasNextTimeRange()) { + curTimeRange = timeRangeIterator.nextTimeRange(); + for (Aggregator aggregator : aggregators) { + aggregator.updateTimeRange(curTimeRange); + } } - // 2. Calculate aggregation result based on current time window + // 1. Calculate aggregation result based on current time window + boolean canCallNext = true; while (!calcFromTsBlock(cachedTsBlock, curTimeRange)) { - if (child.hasNext()) { + cachedTsBlock = null; + // child.next can only be invoked once + if (child.hasNext() && canCallNext) { cachedTsBlock = child.next(); + canCallNext = false; + // if child still has next but can't be invoked now + } else if (child.hasNext()) { + return null; } else { - cachedTsBlock = null; break; } } - // 3. Update result using aggregators + // 2. Update result using aggregators + curTimeRange = null; return updateResultTsBlockFromAggregators(tsBlockBuilder, aggregators, timeRangeIterator); } diff --git a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/SlidingWindowAggregationOperatorTest.java b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/SlidingWindowAggregationOperatorTest.java index ee0deb4df5..b4f849c7b6 100644 --- a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/SlidingWindowAggregationOperatorTest.java +++ b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/SlidingWindowAggregationOperatorTest.java @@ -60,7 +60,6 @@ import java.util.List; import java.util.concurrent.ExecutorService; import java.util.stream.Collectors; -import static org.apache.iotdb.db.constant.TestConstant.count; import static org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext.createFragmentInstanceContext; public class SlidingWindowAggregationOperatorTest { @@ -145,6 +144,9 @@ public class SlidingWindowAggregationOperatorTest { int count = 0; while (slidingWindowAggregationOperator1.hasNext()) { TsBlock resultTsBlock = slidingWindowAggregationOperator1.next(); + if (resultTsBlock == null) { + continue; + } Assert.assertEquals(rootAggregationTypes.size(), resultTsBlock.getValueColumnCount()); Assert.assertEquals(retArray[count], getResultString(resultTsBlock)); count++; @@ -155,6 +157,9 @@ public class SlidingWindowAggregationOperatorTest { initSlidingWindowAggregationOperator(false); while (slidingWindowAggregationOperator2.hasNext()) { TsBlock resultTsBlock = slidingWindowAggregationOperator2.next(); + if (resultTsBlock == null) { + continue; + } Assert.assertEquals(rootAggregationTypes.size(), resultTsBlock.getValueColumnCount()); Assert.assertEquals(retArray[count - 1], getResultString(resultTsBlock)); count--;
