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--;

Reply via email to