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

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

commit 1ac17c5bc218b2d59cedf7660e72da703fac7ba2
Author: Steve Yurong Su <[email protected]>
AuthorDate: Sat Nov 27 22:07:38 2021 +0800

    UDTFFragmentDataSet
---
 .../query/dataset/udf/UDTFAlignByTimeDataSet.java  | 11 ++++++--
 .../iotdb/db/query/dataset/udf/UDTFDataSet.java    | 17 +++++++++--
 .../db/query/dataset/udf/UDTFFragmentDataSet.java  | 33 ++++++++++++++++++++++
 .../db/query/dataset/udf/UDTFNonAlignDataSet.java  |  4 +--
 4 files changed, 58 insertions(+), 7 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFAlignByTimeDataSet.java
 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFAlignByTimeDataSet.java
index 9090538..6484870 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFAlignByTimeDataSet.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFAlignByTimeDataSet.java
@@ -46,7 +46,7 @@ public class UDTFAlignByTimeDataSet extends UDTFDataSet 
implements DirectAlignBy
 
   protected TimeSelector timeHeap;
 
-  /** execute with value filter */
+  /** with value filter */
   public UDTFAlignByTimeDataSet(
       QueryContext context,
       UDTFPlan udtfPlan,
@@ -65,7 +65,7 @@ public class UDTFAlignByTimeDataSet extends UDTFDataSet 
implements DirectAlignBy
     initTimeHeap();
   }
 
-  /** execute without value filter */
+  /** without value filter */
   public UDTFAlignByTimeDataSet(
       QueryContext context, UDTFPlan udtfPlan, List<ManagedSeriesReader> 
readersOfSelectedSeries)
       throws QueryProcessException, IOException, InterruptedException {
@@ -78,6 +78,13 @@ public class UDTFAlignByTimeDataSet extends UDTFDataSet 
implements DirectAlignBy
     initTimeHeap();
   }
 
+  /** for data set fragment */
+  protected UDTFAlignByTimeDataSet(LayerPointReader[] transformers)
+      throws QueryProcessException, IOException {
+    super(transformers);
+    initTimeHeap();
+  }
+
   protected void initTimeHeap() throws IOException, QueryProcessException {
     timeHeap = new TimeSelector(transformers.length << 1, true);
     for (LayerPointReader reader : transformers) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFDataSet.java 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFDataSet.java
index df6e5e4..8f791df 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFDataSet.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFDataSet.java
@@ -39,6 +39,7 @@ import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
 
+/** Used to construct UDTF transformers. */
 public abstract class UDTFDataSet extends QueryDataSet {
 
   protected static final float UDF_READER_MEMORY_BUDGET_IN_MB =
@@ -54,7 +55,7 @@ public abstract class UDTFDataSet extends QueryDataSet {
 
   protected LayerPointReader[] transformers;
 
-  /** execute with value filters */
+  /** with value filters */
   protected UDTFDataSet(
       QueryContext queryContext,
       UDTFPlan udtfPlan,
@@ -80,7 +81,7 @@ public abstract class UDTFDataSet extends QueryDataSet {
     initTransformers();
   }
 
-  /** execute without value filters */
+  /** without value filters */
   protected UDTFDataSet(
       QueryContext queryContext,
       UDTFPlan udtfPlan,
@@ -102,7 +103,7 @@ public abstract class UDTFDataSet extends QueryDataSet {
     initTransformers();
   }
 
-  protected void initTransformers() throws QueryProcessException, IOException {
+  private void initTransformers() throws QueryProcessException, IOException {
     UDFRegistrationService.getInstance().acquireRegistrationLock();
     // This statement must be surrounded by the registration lock.
     UDFClassLoaderManager.getInstance().initializeUDFQuery(queryId);
@@ -123,6 +124,16 @@ public abstract class UDTFDataSet extends QueryDataSet {
     }
   }
 
+  /** for data set fragment */
+  protected UDTFDataSet(LayerPointReader[] transformers) {
+    // The following 3 fields are useless because they are recorded in their 
parent data set.
+    queryId = -1;
+    udtfPlan = null;
+    rawQueryInputLayer = null;
+
+    this.transformers = transformers;
+  }
+
   public void finalizeUDFs(long queryId) {
     udtfPlan.finalizeUDFExecutors(queryId);
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFFragmentDataSet.java
 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFFragmentDataSet.java
new file mode 100644
index 0000000..3e5353c
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFFragmentDataSet.java
@@ -0,0 +1,33 @@
+/*
+ * 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.dataset.udf;
+
+import org.apache.iotdb.db.exception.query.QueryProcessException;
+import org.apache.iotdb.db.query.udf.core.reader.LayerPointReader;
+
+import java.io.IOException;
+
+public class UDTFFragmentDataSet extends UDTFAlignByTimeDataSet {
+
+  protected UDTFFragmentDataSet(LayerPointReader[] transformers)
+      throws QueryProcessException, IOException {
+    super(transformers);
+  }
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFNonAlignDataSet.java
 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFNonAlignDataSet.java
index d13b3ed..999e37b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFNonAlignDataSet.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/dataset/udf/UDTFNonAlignDataSet.java
@@ -51,7 +51,7 @@ public class UDTFNonAlignDataSet extends UDTFDataSet 
implements DirectNonAlignDa
   protected int[] alreadyReturnedRowNumArray;
   protected int[] offsetArray;
 
-  /** execute with value filter */
+  /** with value filter */
   public UDTFNonAlignDataSet(
       QueryContext context,
       UDTFPlan udtfPlan,
@@ -70,7 +70,7 @@ public class UDTFNonAlignDataSet extends UDTFDataSet 
implements DirectNonAlignDa
     isInitialized = false;
   }
 
-  /** execute without value filter */
+  /** without value filter */
   public UDTFNonAlignDataSet(
       QueryContext context, UDTFPlan udtfPlan, List<ManagedSeriesReader> 
readersOfSelectedSeries)
       throws QueryProcessException, IOException, InterruptedException {

Reply via email to