This is an automated email from the ASF dual-hosted git repository.

rong pushed a commit to branch nested-operations
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 9731a895152d26da2acf561d23f9e26c064b753d
Author: Steve Yurong Su <[email protected]>
AuthorDate: Wed Sep 1 12:34:28 2021 +0800

    IntermediateLayer
---
 .../db/query/dataset/UDTFAlignByTimeDataSet.java   |   2 +
 .../apache/iotdb/db/query/dataset/UDTFDataSet.java |   2 +-
 .../db/query/dataset/UDTFNonAlignDataSet.java      |   1 +
 .../db/query/udf/core/builder/DAGBuilder.java      |   5 +
 .../udf/core/builder/LayerPointReaderBuilder.java  |  22 +++
 .../udf/core/{input => layer}/InputLayer.java      |   4 +-
 .../db/query/udf/core/layer/IntermediateLayer.java | 164 +++++++++++++++++++++
 .../udf/core/{input => layer}/SafetyLine.java      |   2 +-
 .../udf/core/transformer/UDFQueryTransformer.java  |   2 +-
 .../tv/ElasticSerializableTVList.java              |   5 +-
 .../ElasticSerializableTVListTest.java             |   2 +-
 11 files changed, 203 insertions(+), 8 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFAlignByTimeDataSet.java
 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFAlignByTimeDataSet.java
index 41ea61a..6c065c2 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFAlignByTimeDataSet.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFAlignByTimeDataSet.java
@@ -192,6 +192,7 @@ public class UDTFAlignByTimeDataSet extends UDTFDataSet 
implements DirectAlignBy
         --rowOffset;
       }
 
+      // todo: control upper bound here
       inputLayer.updateRowRecordListEvictionUpperBound();
     }
 
@@ -294,6 +295,7 @@ public class UDTFAlignByTimeDataSet extends UDTFDataSet 
implements DirectAlignBy
       throw new IOException(e.getMessage());
     }
 
+    // todo: control upper bound here
     inputLayer.updateRowRecordListEvictionUpperBound();
 
     return rowRecord;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFDataSet.java 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFDataSet.java
index b7741ef..0f9ef3d 100644
--- a/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFDataSet.java
+++ b/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFDataSet.java
@@ -36,7 +36,7 @@ import 
org.apache.iotdb.db.query.reader.series.IReaderByTimestamp;
 import org.apache.iotdb.db.query.reader.series.ManagedSeriesReader;
 import org.apache.iotdb.db.query.udf.api.customizer.strategy.AccessStrategy;
 import org.apache.iotdb.db.query.udf.core.executor.UDTFExecutor;
-import org.apache.iotdb.db.query.udf.core.input.InputLayer;
+import org.apache.iotdb.db.query.udf.core.layer.InputLayer;
 import org.apache.iotdb.db.query.udf.core.reader.LayerPointReader;
 import 
org.apache.iotdb.db.query.udf.core.transformer.ArithmeticAdditionTransformer;
 import 
org.apache.iotdb.db.query.udf.core.transformer.ArithmeticDivisionTransformer;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFNonAlignDataSet.java
 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFNonAlignDataSet.java
index 7f90cae..25dd871 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFNonAlignDataSet.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/UDTFNonAlignDataSet.java
@@ -113,6 +113,7 @@ public class UDTFNonAlignDataSet extends UDTFDataSet 
implements DirectNonAlignDa
       valueBufferList.add(timeValueByteBufferPair.right);
     }
 
+    // todo: control upper bound here
     inputLayer.updateRowRecordListEvictionUpperBound();
 
     tsQueryNonAlignDataSet.setTimeList(timeBufferList);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/builder/DAGBuilder.java
 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/builder/DAGBuilder.java
index eaaaa14..d8dc64e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/builder/DAGBuilder.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/builder/DAGBuilder.java
@@ -31,6 +31,8 @@ import java.util.Map;
 
 public class DAGBuilder {
 
+  private final UDTFPlan udtfPlan;
+
   // input
   private final List<Expression> resultColumnExpressions;
   // output
@@ -43,12 +45,15 @@ public class DAGBuilder {
   private final Map<Expression, TransformerBuilder> 
expressionTransformerBuilderMap;
 
   public DAGBuilder(UDTFPlan udtfPlan) {
+    this.udtfPlan = udtfPlan;
     resultColumnExpressions = new ArrayList<>();
     for (ResultColumn resultColumn : udtfPlan.getResultColumns()) {
       resultColumnExpressions.add(resultColumn.getExpression());
     }
     resultColumnTransformers = new Transformer[resultColumnExpressions.size()];
     expressionTransformerBuilderMap = new HashMap<>();
+
+    build();
   }
 
   public void build() {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/builder/LayerPointReaderBuilder.java
 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/builder/LayerPointReaderBuilder.java
new file mode 100644
index 0000000..ed96576
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/builder/LayerPointReaderBuilder.java
@@ -0,0 +1,22 @@
+/*
+ * 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.query.udf.core.builder;
+
+public class LayerPointReaderBuilder {}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/input/InputLayer.java 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/InputLayer.java
similarity index 99%
rename from 
server/src/main/java/org/apache/iotdb/db/query/udf/core/input/InputLayer.java
rename to 
server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/InputLayer.java
index de08f96..f16a210 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/input/InputLayer.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/InputLayer.java
@@ -17,7 +17,7 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.query.udf.core.input;
+package org.apache.iotdb.db.query.udf.core.layer;
 
 import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.metadata.PartialPath;
@@ -33,7 +33,7 @@ import 
org.apache.iotdb.db.query.udf.api.customizer.strategy.SlidingSizeWindowAc
 import 
org.apache.iotdb.db.query.udf.api.customizer.strategy.SlidingTimeWindowAccessStrategy;
 import org.apache.iotdb.db.query.udf.core.access.RowImpl;
 import org.apache.iotdb.db.query.udf.core.access.RowWindowImpl;
-import org.apache.iotdb.db.query.udf.core.input.SafetyLine.SafetyPile;
+import org.apache.iotdb.db.query.udf.core.layer.SafetyLine.SafetyPile;
 import org.apache.iotdb.db.query.udf.core.reader.LayerPointReader;
 import org.apache.iotdb.db.query.udf.core.reader.LayerRowReader;
 import org.apache.iotdb.db.query.udf.core.reader.LayerRowWindowReader;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/IntermediateLayer.java
 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/IntermediateLayer.java
new file mode 100644
index 0000000..f4ccd03
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/IntermediateLayer.java
@@ -0,0 +1,164 @@
+/*
+ * 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.query.udf.core.layer;
+
+import org.apache.iotdb.db.exception.query.QueryProcessException;
+import org.apache.iotdb.db.query.udf.core.layer.SafetyLine.SafetyPile;
+import org.apache.iotdb.db.query.udf.core.reader.LayerPointReader;
+import 
org.apache.iotdb.db.query.udf.datastructure.tv.ElasticSerializableTVList;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.utils.Binary;
+
+import java.io.IOException;
+
+public class IntermediateLayer {
+
+  private static final int CACHE_BLOCK_SIZE = 1;
+
+  private final TSDataType dataType;
+  private final LayerPointReader parentLayerPointReader;
+  private final ElasticSerializableTVList tvList;
+  private final SafetyLine safetyLine;
+
+  public IntermediateLayer(
+      LayerPointReader parentLayerPointReader, long queryId, float 
memoryBudgetInMB)
+      throws QueryProcessException {
+    this.parentLayerPointReader = parentLayerPointReader;
+    dataType = parentLayerPointReader.getDataType();
+    tvList =
+        ElasticSerializableTVList.newElasticSerializableTVList(
+            dataType, queryId, memoryBudgetInMB, CACHE_BLOCK_SIZE);
+    safetyLine = new SafetyLine();
+  }
+
+  public LayerPointReader constructPointReader() {
+
+    return new LayerPointReader() {
+
+      private final SafetyPile safetyPile = safetyLine.addSafetyPile();
+
+      private boolean hasCached = false;
+      private int currentPointIndex = -1;
+
+      @Override
+      public boolean next() throws QueryProcessException, IOException {
+        if (hasCached) {
+          return true;
+        }
+
+        if (currentPointIndex < tvList.size() - 1) {
+          ++currentPointIndex;
+          hasCached = true;
+        }
+
+        // tvList.size() - 1 <= currentPointIndex
+        if (!hasCached && parentLayerPointReader.next()) {
+          cachePoint();
+          parentLayerPointReader.readyForNext();
+
+          ++currentPointIndex;
+          hasCached = true;
+        }
+
+        return hasCached;
+      }
+
+      private void cachePoint() throws IOException, QueryProcessException {
+        switch (dataType) {
+          case INT32:
+            tvList.putInt(
+                parentLayerPointReader.currentTime(), 
parentLayerPointReader.currentInt());
+            break;
+          case INT64:
+            tvList.putLong(
+                parentLayerPointReader.currentTime(), 
parentLayerPointReader.currentLong());
+            break;
+          case FLOAT:
+            tvList.putFloat(
+                parentLayerPointReader.currentTime(), 
parentLayerPointReader.currentFloat());
+            break;
+          case DOUBLE:
+            tvList.putDouble(
+                parentLayerPointReader.currentTime(), 
parentLayerPointReader.currentDouble());
+            break;
+          case BOOLEAN:
+            tvList.putBoolean(
+                parentLayerPointReader.currentTime(), 
parentLayerPointReader.currentBoolean());
+            break;
+          case TEXT:
+            tvList.putBinary(
+                parentLayerPointReader.currentTime(), 
parentLayerPointReader.currentBinary());
+            break;
+          default:
+            throw new UnsupportedOperationException(dataType.name());
+        }
+      }
+
+      @Override
+      public void readyForNext() {
+        hasCached = false;
+
+        safetyPile.moveForwardTo(currentPointIndex + 1);
+        // todo: reduce the update rate
+        tvList.setEvictionUpperBound(safetyLine.getSafetyLine());
+      }
+
+      @Override
+      public TSDataType getDataType() {
+        return dataType;
+      }
+
+      @Override
+      public long currentTime() throws IOException {
+        return tvList.getTime(currentPointIndex);
+      }
+
+      @Override
+      public int currentInt() throws IOException {
+        return tvList.getInt(currentPointIndex);
+      }
+
+      @Override
+      public long currentLong() throws IOException {
+        return tvList.getLong(currentPointIndex);
+      }
+
+      @Override
+      public float currentFloat() throws IOException {
+        return tvList.getFloat(currentPointIndex);
+      }
+
+      @Override
+      public double currentDouble() throws IOException {
+        return tvList.getDouble(currentPointIndex);
+      }
+
+      @Override
+      public boolean currentBoolean() throws IOException {
+        return tvList.getBoolean(currentPointIndex);
+      }
+
+      @Override
+      public Binary currentBinary() throws IOException {
+        return tvList.getBinary(currentPointIndex);
+      }
+    };
+  }
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/input/SafetyLine.java 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/SafetyLine.java
similarity index 97%
rename from 
server/src/main/java/org/apache/iotdb/db/query/udf/core/input/SafetyLine.java
rename to 
server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/SafetyLine.java
index 7f21c90..35be04d 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/input/SafetyLine.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/layer/SafetyLine.java
@@ -17,7 +17,7 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.query.udf.core.input;
+package org.apache.iotdb.db.query.udf.core.layer;
 
 public class SafetyLine {
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/transformer/UDFQueryTransformer.java
 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/transformer/UDFQueryTransformer.java
index 7bb5b60..5254a50 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/udf/core/transformer/UDFQueryTransformer.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/udf/core/transformer/UDFQueryTransformer.java
@@ -39,7 +39,7 @@ public abstract class UDFQueryTransformer extends Transformer 
{
   protected UDFQueryTransformer(UDTFExecutor executor) {
     this.executor = executor;
     udfOutputDataType = executor.getConfigurations().getOutputDataType();
-    udfOutput = executor.getCollector().getPointReaderUsingEvictionStrategy();
+    udfOutput = 
executor.getCollector().constructPointReaderUsingTrivialEvictionStrategy();
     terminated = false;
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/udf/datastructure/tv/ElasticSerializableTVList.java
 
b/server/src/main/java/org/apache/iotdb/db/query/udf/datastructure/tv/ElasticSerializableTVList.java
index 6dd0a9f..8d36911 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/udf/datastructure/tv/ElasticSerializableTVList.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/udf/datastructure/tv/ElasticSerializableTVList.java
@@ -224,7 +224,8 @@ public class ElasticSerializableTVList implements 
PointCollector {
     }
   }
 
-  public LayerPointReader getPointReaderUsingEvictionStrategy() {
+  // todo: remove it
+  public LayerPointReader constructPointReaderUsingTrivialEvictionStrategy() {
 
     return new LayerPointReader() {
 
@@ -290,7 +291,7 @@ public class ElasticSerializableTVList implements 
PointCollector {
    * @param evictionUpperBound the index of the first element that cannot be 
evicted. in other
    *     words, elements whose index are <b>less than</b> the 
evictionUpperBound can be evicted.
    */
-  private void setEvictionUpperBound(int evictionUpperBound) {
+  public void setEvictionUpperBound(int evictionUpperBound) {
     this.evictionUpperBound = evictionUpperBound;
   }
 
diff --git 
a/server/src/test/java/org/apache/iotdb/db/query/udf/datastructure/ElasticSerializableTVListTest.java
 
b/server/src/test/java/org/apache/iotdb/db/query/udf/datastructure/ElasticSerializableTVListTest.java
index db287b4..280e340 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/query/udf/datastructure/ElasticSerializableTVListTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/query/udf/datastructure/ElasticSerializableTVListTest.java
@@ -202,7 +202,7 @@ public class ElasticSerializableTVListTest extends 
SerializableListTest {
                 generateRandomString(
                     byteLengthMin + random.nextInt(byteLengthMax - 
byteLengthMin))));
       }
-      LayerPointReader reader = tvList.getPointReaderUsingEvictionStrategy();
+      LayerPointReader reader = 
tvList.constructPointReaderUsingTrivialEvictionStrategy();
       while (reader.next()) {
         int length = reader.currentBinary().getLength();
         assertTrue(byteLengthMin <= length && length < byteLengthMax);

Reply via email to