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 6b95392327a03576f45718ae8c5a685bbef153fe
Author: Alima777 <[email protected]>
AuthorDate: Tue May 31 16:50:58 2022 +0800

    modify RawDataAggregationOperator control next logic
---
 .../process/RawDataAggregationOperator.java        | 28 +++++++++++++++-------
 .../operator/RawDataAggregationOperatorTest.java   | 12 ++++++++++
 2 files changed, 31 insertions(+), 9 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 3753635e7f..683d32fbac 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
@@ -94,30 +94,40 @@ public class RawDataAggregationOperator implements 
ProcessOperator {
 
   @Override
   public TsBlock next() {
-    // 1. Clear previous aggregation result
-    curTimeRange = timeRangeIterator.nextTimeRange();
-    for (Aggregator aggregator : aggregators) {
-      aggregator.reset();
-      aggregator.updateTimeRange(curTimeRange);
+    // Move to next timeRange
+    if (curTimeRange == null && timeRangeIterator.hasNextTimeRange()) {
+      curTimeRange = timeRangeIterator.nextTimeRange();
+      for (Aggregator aggregator : aggregators) {
+        aggregator.reset();
+        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 (!calcFromCacheData(curTimeRange)) {
-      if (child.hasNext()) {
+      // child.next can only be invoked once
+      if (child.hasNext() && canCallNext) {
+        canCallNext = false;
         preCachedData = child.next();
+        // if child still has next but can't be invoked now
+      } else if (child.hasNext()) {
+        preCachedData = null;
+        return null;
       } else {
         break;
       }
     }
 
-    // 3. Update result using aggregators
+    // 2. Update result using aggregators
+    curTimeRange = null;
     return AggregationOperator.updateResultTsBlockFromAggregators(
         tsBlockBuilder, aggregators, timeRangeIterator);
   }
 
   @Override
   public boolean hasNext() {
-    return timeRangeIterator.hasNextTimeRange();
+    return curTimeRange != null || timeRangeIterator.hasNextTimeRange();
   }
 
   @Override
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/RawDataAggregationOperatorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/RawDataAggregationOperatorTest.java
index daed697a13..46b85b85f7 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/RawDataAggregationOperatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/RawDataAggregationOperatorTest.java
@@ -122,6 +122,9 @@ public class RawDataAggregationOperatorTest {
     int count = 0;
     while (rawDataAggregationOperator.hasNext()) {
       TsBlock resultTsBlock = rawDataAggregationOperator.next();
+      if (resultTsBlock == null) {
+        continue;
+      }
       for (int i = 0; i < 2; i++) {
         assertEquals(500, resultTsBlock.getColumn(6 * i).getLong(0));
         assertEquals(6524750.0, resultTsBlock.getColumn(6 * i + 
1).getDouble(0), 0.0001);
@@ -172,6 +175,9 @@ public class RawDataAggregationOperatorTest {
     int count = 0;
     while (rawDataAggregationOperator.hasNext()) {
       TsBlock resultTsBlock = rawDataAggregationOperator.next();
+      if (resultTsBlock == null) {
+        continue;
+      }
       for (int i = 0; i < 2; i++) {
         assertEquals(13049.5, resultTsBlock.getColumn(i).getDouble(0), 0.001);
       }
@@ -220,6 +226,9 @@ public class RawDataAggregationOperatorTest {
     int count = 0;
     while (rawDataAggregationOperator.hasNext()) {
       TsBlock resultTsBlock = rawDataAggregationOperator.next();
+      if (resultTsBlock == null) {
+        continue;
+      }
       assertEquals(100 * count, resultTsBlock.getTimeColumn().getLong(0));
       for (int i = 0; i < 2; i++) {
         assertEquals(result[0][count], resultTsBlock.getColumn(6 * 
i).getLong(0));
@@ -269,6 +278,9 @@ public class RawDataAggregationOperatorTest {
     int count = 0;
     while (rawDataAggregationOperator.hasNext()) {
       TsBlock resultTsBlock = rawDataAggregationOperator.next();
+      if (resultTsBlock == null) {
+        continue;
+      }
       assertEquals(100 * count, resultTsBlock.getTimeColumn().getLong(0));
       for (int i = 0; i < 2; i++) {
         assertEquals(result[0][count], 
resultTsBlock.getColumn(i).getDouble(0), 0.001);

Reply via email to