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 {
