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