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