This is an automated email from the ASF dual-hosted git repository. jt2594838 pushed a commit to branch remove_swtich_type in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 49a2e7a1d12a5430c6c6c98683e28ccac8582f2f Author: Tian Jiang <[email protected]> AuthorDate: Fri Aug 28 11:57:33 2026 +0800 multiple refactors --- .../iotdb/commons/udf/builtin/TypeServices.java | 414 +++++++++++++++++++++ .../apache/iotdb/commons/udf/builtin/UDTFAbs.java | 106 +----- .../iotdb/commons/udf/builtin/UDTFBottomK.java | 105 +++--- .../commons/udf/builtin/UDTFChangePoints.java | 178 ++++----- .../iotdb/commons/udf/builtin/UDTFConst.java | 267 +++++-------- .../udf/builtin/UDTFContinuouslySatisfy.java | 79 +--- .../udf/builtin/UDTFEqualSizeBucketAggSample.java | 53 ++- .../udf/builtin/UDTFEqualSizeBucketM4Sample.java | 39 +- .../builtin/UDTFEqualSizeBucketOutlierSample.java | 53 ++- .../builtin/UDTFEqualSizeBucketRandomSample.java | 40 +- .../apache/iotdb/commons/udf/builtin/UDTFM4.java | 36 +- .../iotdb/commons/udf/builtin/UDTFSelectK.java | 136 +++---- .../apache/iotdb/commons/udf/builtin/UDTFTopK.java | 65 ++-- 13 files changed, 791 insertions(+), 780 deletions(-) diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/TypeServices.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/TypeServices.java index 2957e953ea2..b4c8333a3fe 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/TypeServices.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/TypeServices.java @@ -20,11 +20,14 @@ package org.apache.iotdb.commons.udf.builtin; import org.apache.iotdb.commons.udf.utils.UDFDataTypeTransformer; import org.apache.iotdb.udf.api.access.Row; +import org.apache.iotdb.udf.api.access.RowWindow; import org.apache.iotdb.udf.api.collector.PointCollector; +import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; import org.apache.iotdb.udf.api.type.Type; import org.apache.tsfile.block.column.Column; +import org.apache.tsfile.block.column.ColumnBuilder; import org.apache.tsfile.read.common.type.service.TypeService; import java.io.IOException; @@ -203,6 +206,265 @@ final class TypeServices { }; }; + static final TypeService<NumericRowCollector> NUMERIC_ROW_COLLECTOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32 -> (row, collector) -> collector.putInt(row.getTime(), row.getInt(0)); + case INT64 -> (row, collector) -> collector.putLong(row.getTime(), row.getLong(0)); + case FLOAT -> (row, collector) -> collector.putFloat(row.getTime(), row.getFloat(0)); + case DOUBLE -> (row, collector) -> collector.putDouble(row.getTime(), row.getDouble(0)); + case BOOLEAN, DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (row, collector) -> { + throw invalidNumericDataType(type); + }; + }; + + static final TypeService<AbsRowCollector> ABS_ROW_COLLECTOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32 -> + (row, collector) -> collector.putInt(row.getTime(), Math.abs(row.getInt(0))); + case INT64 -> + (row, collector) -> collector.putLong(row.getTime(), Math.abs(row.getLong(0))); + case FLOAT -> + (row, collector) -> collector.putFloat(row.getTime(), Math.abs(row.getFloat(0))); + case DOUBLE -> + (row, collector) -> collector.putDouble(row.getTime(), Math.abs(row.getDouble(0))); + case BOOLEAN, DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (row, collector) -> { + throw invalidNumericDataType(type); + }; + }; + + static final TypeService<AbsRowMapper> ABS_ROW_MAPPER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32 -> row -> Math.abs(row.getInt(0)); + case INT64 -> row -> Math.abs(row.getLong(0)); + case FLOAT -> row -> Math.abs(row.getFloat(0)); + case DOUBLE -> row -> Math.abs(row.getDouble(0)); + case BOOLEAN, DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + row -> { + throw invalidNumericDataType(type); + }; + }; + + static final TypeService<AbsColumnTransformer> ABS_COLUMN_TRANSFORMER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32 -> UDTFAbs::transformInt; + case INT64 -> UDTFAbs::transformLong; + case FLOAT -> UDTFAbs::transformFloat; + case DOUBLE -> UDTFAbs::transformDouble; + case BOOLEAN, DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (target, columns, builder) -> { + throw invalidNumericDataType(type); + }; + }; + + static final TypeService<NumericWindowTransformer<UDTFM4>> M4_WINDOW_TRANSFORMER_SERVICE = + numericWindowTransformerService( + UDTFM4::transformInt, + UDTFM4::transformLong, + UDTFM4::transformFloat, + UDTFM4::transformDouble); + + static final TypeService<NumericWindowTransformer<UDTFEqualSizeBucketM4Sample>> + BUCKET_M4_WINDOW_TRANSFORMER_SERVICE = + numericWindowTransformerService( + UDTFEqualSizeBucketM4Sample::transformInt, + UDTFEqualSizeBucketM4Sample::transformLong, + UDTFEqualSizeBucketM4Sample::transformFloat, + UDTFEqualSizeBucketM4Sample::transformDouble); + + static final TypeService<NumericWindowTransformer<UDTFEqualSizeBucketAggSample>> + BUCKET_AGG_WINDOW_TRANSFORMER_SERVICE = + numericWindowTransformerService( + UDTFEqualSizeBucketAggSample::aggregateInt, + UDTFEqualSizeBucketAggSample::aggregateLong, + UDTFEqualSizeBucketAggSample::aggregateFloat, + UDTFEqualSizeBucketAggSample::aggregateDouble); + + static final TypeService<NumericWindowTransformer<UDTFEqualSizeBucketOutlierSample>> + BUCKET_OUTLIER_WINDOW_TRANSFORMER_SERVICE = + numericWindowTransformerService( + UDTFEqualSizeBucketOutlierSample::outlierSampleInt, + UDTFEqualSizeBucketOutlierSample::outlierSampleLong, + UDTFEqualSizeBucketOutlierSample::outlierSampleFloat, + UDTFEqualSizeBucketOutlierSample::outlierSampleDouble); + + static final TypeService<ChangePointProcessor> CHANGE_POINT_PROCESSOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case BOOLEAN -> UDTFChangePoints::transformBoolean; + case INT32 -> UDTFChangePoints::transformInt; + case INT64 -> UDTFChangePoints::transformLong; + case FLOAT -> UDTFChangePoints::transformFloat; + case DOUBLE -> UDTFChangePoints::transformDouble; + case TEXT -> UDTFChangePoints::transformString; + case DATE, TIMESTAMP, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (target, row, collector) -> {}; + }; + + static final TypeService<SelectKRowTransformer> SELECT_K_ROW_TRANSFORMER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, DATE -> (target, row) -> target.transformInt(row.getTime(), row.getInt(0)); + case INT64, TIMESTAMP -> + (target, row) -> target.transformLong(row.getTime(), row.getLong(0)); + case FLOAT -> (target, row) -> target.transformFloat(row.getTime(), row.getFloat(0)); + case DOUBLE -> (target, row) -> target.transformDouble(row.getTime(), row.getDouble(0)); + case TEXT, STRING -> + (target, row) -> target.transformString(row.getTime(), row.getString(0)); + case BOOLEAN, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (target, row) -> { + throw invalidSelectKDataType(type); + }; + }; + + static final TypeService<SelectKTerminator> SELECT_K_TERMINATOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, DATE -> UDTFSelectK::terminateInt; + case INT64, TIMESTAMP -> UDTFSelectK::terminateLong; + case FLOAT -> UDTFSelectK::terminateFloat; + case DOUBLE -> UDTFSelectK::terminateDouble; + case TEXT, STRING -> UDTFSelectK::terminateString; + case BOOLEAN, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (target, collector) -> { + throw invalidSelectKDataType(type); + }; + }; + + static final TypeService<TopKQueueConstructor> TOP_K_QUEUE_CONSTRUCTOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, DATE -> UDTFTopK::initializeIntQueue; + case INT64, TIMESTAMP -> UDTFTopK::initializeLongQueue; + case FLOAT -> UDTFTopK::initializeFloatQueue; + case DOUBLE -> UDTFTopK::initializeDoubleQueue; + case TEXT, STRING -> UDTFTopK::initializeStringQueue; + case BOOLEAN, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + target -> { + throw invalidSelectKDataType(type); + }; + }; + + static final TypeService<BottomKQueueConstructor> BOTTOM_K_QUEUE_CONSTRUCTOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, DATE -> UDTFBottomK::initializeIntQueue; + case INT64, TIMESTAMP -> UDTFBottomK::initializeLongQueue; + case FLOAT -> UDTFBottomK::initializeFloatQueue; + case DOUBLE -> UDTFBottomK::initializeDoubleQueue; + case TEXT, STRING -> UDTFBottomK::initializeStringQueue; + case BOOLEAN, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + target -> { + throw invalidSelectKDataType(type); + }; + }; + + static final TypeService<ConstantParser> CONSTANT_PARSER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32 -> UDTFConst::parseInt; + case DATE -> UDTFConst::parseDate; + case INT64, TIMESTAMP -> UDTFConst::parseLong; + case FLOAT -> UDTFConst::parseFloat; + case DOUBLE -> UDTFConst::parseDouble; + case BOOLEAN -> UDTFConst::parseBoolean; + case TEXT, STRING -> UDTFConst::parseText; + case BLOB, OBJECT -> UDTFConst::parseBlob; + case ROW, UNKNOWN, VECTOR -> + (target, parameters) -> { + throw new UnsupportedOperationException(); + }; + }; + + static final TypeService<ConstantRowCollector> CONSTANT_ROW_COLLECTOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, DATE -> + (target, row, collector) -> collector.putInt(row.getTime(), target.intValue()); + case INT64, TIMESTAMP -> + (target, row, collector) -> collector.putLong(row.getTime(), target.longValue()); + case FLOAT -> + (target, row, collector) -> collector.putFloat(row.getTime(), target.floatValue()); + case DOUBLE -> + (target, row, collector) -> + collector.putDouble(row.getTime(), target.doubleValue()); + case BOOLEAN -> + (target, row, collector) -> + collector.putBoolean(row.getTime(), target.booleanValue()); + case TEXT, STRING, BLOB, OBJECT -> + (target, row, collector) -> + collector.putBinary(row.getTime(), target.binaryValue()); + case ROW, UNKNOWN, VECTOR -> + (target, row, collector) -> { + throw new UnsupportedOperationException(); + }; + }; + + static final TypeService<ConstantRowMapper> CONSTANT_ROW_MAPPER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, DATE -> UDTFConst::intValue; + case INT64, TIMESTAMP -> UDTFConst::longValue; + case FLOAT -> UDTFConst::floatValue; + case DOUBLE -> UDTFConst::doubleValue; + case BOOLEAN -> UDTFConst::booleanValue; + case TEXT, STRING, BLOB, OBJECT -> UDTFConst::binaryValue; + case ROW, UNKNOWN, VECTOR -> + target -> { + throw new UnsupportedOperationException(); + }; + }; + + static final TypeService<ConstantColumnValueWriter> CONSTANT_COLUMN_VALUE_WRITER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, DATE -> (target, builder) -> builder.writeInt(target.intValue()); + case INT64, TIMESTAMP -> (target, builder) -> builder.writeLong(target.longValue()); + case FLOAT -> (target, builder) -> builder.writeFloat(target.floatValue()); + case DOUBLE -> (target, builder) -> builder.writeDouble(target.doubleValue()); + case BOOLEAN -> (target, builder) -> builder.writeBoolean(target.booleanValue()); + case TEXT, STRING, BLOB, OBJECT -> + (target, builder) -> builder.writeBinary(target.tsFileBinaryValue()); + case ROW, UNKNOWN, VECTOR -> + (target, builder) -> { + throw new UnsupportedOperationException(); + }; + }; + + static final TypeService<ContinuouslySatisfyRowTransformer> + CONTINUOUSLY_SATISFY_ROW_TRANSFORMER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32 -> (target, row) -> target.transformInt(row.getTime(), row.getInt(0)); + case INT64 -> (target, row) -> target.transformLong(row.getTime(), row.getLong(0)); + case FLOAT -> + (target, row) -> target.transformFloat(row.getTime(), row.getFloat(0)); + case DOUBLE -> + (target, row) -> target.transformDouble(row.getTime(), row.getDouble(0)); + case BOOLEAN -> + (target, row) -> target.transformBoolean(row.getTime(), row.getBoolean(0)); + case DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (target, row) -> { + throw invalidContinuouslySatisfyDataType(type); + }; + }; + + static final TypeService<ContinuouslySatisfyTerminator> CONTINUOUSLY_SATISFY_TERMINATOR_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, INT64, FLOAT, DOUBLE, BOOLEAN -> + UDTFContinuouslySatisfy::terminateSupportedType; + case DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (target, collector) -> { + throw invalidContinuouslySatisfyDataType(type); + }; + }; + static { VALUE_TREND_READER_SERVICE.check(); VALUE_DIFFERENCE_OPERATOR_SERVICE.check(); @@ -211,6 +473,25 @@ final class TypeServices { NON_NEGATIVE_DERIVATIVE_OPERATOR_SERVICE.check(); NUMERIC_ROW_READER_SERVICE.check(); NUMERIC_COLUMN_READER_SERVICE.check(); + NUMERIC_ROW_COLLECTOR_SERVICE.check(); + ABS_ROW_COLLECTOR_SERVICE.check(); + ABS_ROW_MAPPER_SERVICE.check(); + ABS_COLUMN_TRANSFORMER_SERVICE.check(); + M4_WINDOW_TRANSFORMER_SERVICE.check(); + BUCKET_M4_WINDOW_TRANSFORMER_SERVICE.check(); + BUCKET_AGG_WINDOW_TRANSFORMER_SERVICE.check(); + BUCKET_OUTLIER_WINDOW_TRANSFORMER_SERVICE.check(); + CHANGE_POINT_PROCESSOR_SERVICE.check(); + SELECT_K_ROW_TRANSFORMER_SERVICE.check(); + SELECT_K_TERMINATOR_SERVICE.check(); + TOP_K_QUEUE_CONSTRUCTOR_SERVICE.check(); + BOTTOM_K_QUEUE_CONSTRUCTOR_SERVICE.check(); + CONSTANT_PARSER_SERVICE.check(); + CONSTANT_ROW_COLLECTOR_SERVICE.check(); + CONSTANT_ROW_MAPPER_SERVICE.check(); + CONSTANT_COLUMN_VALUE_WRITER_SERVICE.check(); + CONTINUOUSLY_SATISFY_ROW_TRANSFORMER_SERVICE.check(); + CONTINUOUSLY_SATISFY_TERMINATOR_SERVICE.check(); } private TypeServices() {} @@ -226,6 +507,51 @@ final class TypeServices { Type.DOUBLE); } + private static UDFInputSeriesDataTypeNotValidException invalidSelectKDataType( + org.apache.tsfile.read.common.type.Type type) { + return new UDFInputSeriesDataTypeNotValidException( + 0, + UDFDataTypeTransformer.transformReadTypeToUDFDataType(type), + Type.INT32, + Type.INT64, + Type.FLOAT, + Type.DOUBLE, + Type.TEXT, + Type.DATE, + Type.TIMESTAMP, + Type.STRING); + } + + private static UDFInputSeriesDataTypeNotValidException invalidContinuouslySatisfyDataType( + org.apache.tsfile.read.common.type.Type type) { + return new UDFInputSeriesDataTypeNotValidException( + 0, + UDFDataTypeTransformer.transformReadTypeToUDFDataType(type), + Type.INT32, + Type.INT64, + Type.FLOAT, + Type.DOUBLE); + } + + // Keep each window algorithm's primitive output type while selecting its callback once. + private static <T> TypeService<NumericWindowTransformer<T>> numericWindowTransformerService( + NumericWindowTransformer<T> intTransformer, + NumericWindowTransformer<T> longTransformer, + NumericWindowTransformer<T> floatTransformer, + NumericWindowTransformer<T> doubleTransformer) { + return type -> + switch (type.getTypeEnum()) { + case INT32 -> intTransformer; + case INT64 -> longTransformer; + case FLOAT -> floatTransformer; + case DOUBLE -> doubleTransformer; + case BOOLEAN, DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (target, rowWindow, collector) -> { + throw invalidNumericDataType(type); + }; + }; + } + @FunctionalInterface interface PreviousValueReader { void read(UDTFValueTrend target, Row row) @@ -254,4 +580,92 @@ final class TypeServices { interface NumericColumnReader { double read(Column column, int position) throws UDFInputSeriesDataTypeNotValidException; } + + @FunctionalInterface + interface NumericRowCollector { + void collect(Row row, PointCollector collector) + throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface AbsRowCollector { + void collect(Row row, PointCollector collector) + throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface AbsRowMapper { + Object map(Row row) throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface AbsColumnTransformer { + void transform(UDTFAbs target, Column[] columns, ColumnBuilder builder) + throws UDFInputSeriesDataTypeNotValidException; + } + + @FunctionalInterface + interface NumericWindowTransformer<T> { + void transform(T target, RowWindow rowWindow, PointCollector collector) + throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface ChangePointProcessor { + void transform(UDTFChangePoints target, Row row, PointCollector collector) throws IOException; + } + + @FunctionalInterface + interface SelectKRowTransformer { + void transform(UDTFSelectK target, Row row) + throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface SelectKTerminator { + void terminate(UDTFSelectK target, PointCollector collector) + throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface TopKQueueConstructor { + void construct(UDTFTopK target); + } + + @FunctionalInterface + interface BottomKQueueConstructor { + void construct(UDTFBottomK target); + } + + @FunctionalInterface + interface ConstantParser { + void parse(UDTFConst target, UDFParameters parameters); + } + + @FunctionalInterface + interface ConstantRowCollector { + void collect(UDTFConst target, Row row, PointCollector collector) throws IOException; + } + + @FunctionalInterface + interface ConstantRowMapper { + Object map(UDTFConst target); + } + + @FunctionalInterface + interface ConstantColumnValueWriter { + void write(UDTFConst target, ColumnBuilder builder); + } + + @FunctionalInterface + interface ContinuouslySatisfyRowTransformer { + boolean transform(UDTFContinuouslySatisfy target, Row row) + throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface ContinuouslySatisfyTerminator { + void terminate(UDTFContinuouslySatisfy target, PointCollector collector) + throws UDFInputSeriesDataTypeNotValidException, IOException; + } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFAbs.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFAbs.java index 2af7ee3c83d..06303b8d2a3 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFAbs.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFAbs.java @@ -27,62 +27,36 @@ import org.apache.iotdb.udf.api.collector.PointCollector; import org.apache.iotdb.udf.api.customizer.config.UDTFConfigurations; import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.customizer.strategy.MappableRowByRowAccessStrategy; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; -import org.apache.iotdb.udf.api.type.Type; import org.apache.tsfile.block.column.Column; import org.apache.tsfile.block.column.ColumnBuilder; +import org.apache.tsfile.read.common.type.Type; import java.io.IOException; @SuppressWarnings("java:S2177") public class UDTFAbs extends UDTFMath { + private TypeServices.AbsRowCollector rowCollector; + private TypeServices.AbsRowMapper rowMapper; + private TypeServices.AbsColumnTransformer columnTransformer; + @Override public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) throws MetadataException { dataType = UDFDataTypeTransformer.transformToTsDataType(parameters.getDataType(0)); + Type type = Type.fromTsDataType(dataType); + rowCollector = TypeServices.ABS_ROW_COLLECTOR_SERVICE.call(type); + rowMapper = TypeServices.ABS_ROW_MAPPER_SERVICE.call(type); + columnTransformer = TypeServices.ABS_COLUMN_TRANSFORMER_SERVICE.call(type); configurations .setAccessStrategy(new MappableRowByRowAccessStrategy()) .setOutputDataType(UDFDataTypeTransformer.transformToUDFDataType(dataType)); } @Override - public void transform(Row row, PointCollector collector) - throws UDFInputSeriesDataTypeNotValidException, IOException { - long time = row.getTime(); - switch (dataType) { - case INT32: - collector.putInt(time, Math.abs(row.getInt(0))); - break; - case INT64: - collector.putLong(time, Math.abs(row.getLong(0))); - break; - case FLOAT: - collector.putFloat(time, Math.abs(row.getFloat(0))); - break; - case DOUBLE: - collector.putDouble(time, Math.abs(row.getDouble(0))); - break; - case BLOB: - case OBJECT: - case STRING: - case TIMESTAMP: - case TEXT: - case DATE: - case BOOLEAN: - case VECTOR: - case UNKNOWN: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + public void transform(Row row, PointCollector collector) throws IOException { + rowCollector.collect(row, collector); } @Override @@ -90,66 +64,12 @@ public class UDTFAbs extends UDTFMath { if (row.isNull(0)) { return null; } - switch (dataType) { - case INT32: - return Math.abs(row.getInt(0)); - case INT64: - return Math.abs(row.getLong(0)); - case FLOAT: - return Math.abs(row.getFloat(0)); - case DOUBLE: - return Math.abs(row.getDouble(0)); - case DATE: - case BOOLEAN: - case TEXT: - case TIMESTAMP: - case STRING: - case BLOB: - case OBJECT: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + return rowMapper.map(row); } @Override public void transform(Column[] columns, ColumnBuilder builder) throws Exception { - switch (dataType) { - case INT32: - transformInt(columns, builder); - return; - case INT64: - transformLong(columns, builder); - return; - case FLOAT: - transformFloat(columns, builder); - return; - case DOUBLE: - transformDouble(columns, builder); - return; - case BLOB: - case OBJECT: - case STRING: - case TEXT: - case TIMESTAMP: - case BOOLEAN: - case DATE: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + columnTransformer.transform(this, columns, builder); } @Override diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFBottomK.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFBottomK.java index 045e49b064b..b4702c8046b 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFBottomK.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFBottomK.java @@ -19,10 +19,6 @@ package org.apache.iotdb.commons.udf.builtin; -import org.apache.iotdb.commons.udf.utils.UDFDataTypeTransformer; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; -import org.apache.iotdb.udf.api.type.Type; - import org.apache.tsfile.utils.Pair; import java.util.Comparator; @@ -32,63 +28,50 @@ import java.util.PriorityQueue; public class UDTFBottomK extends UDTFSelectK { @Override - protected void constructPQ() throws UDFInputSeriesDataTypeNotValidException { - switch (dataType) { - case INT32: - case DATE: - intPQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); - break; - case INT64: - case TIMESTAMP: - longPQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); - break; - case FLOAT: - floatPQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); - break; - case DOUBLE: - doublePQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); - break; - case TEXT: - case STRING: - stringPQ = - new PriorityQueue<>( - k, - (pairA, pairB) -> { - final String cs1 = pairB.right; - final String cs2 = pairA.right; - - if (Objects.requireNonNull(cs1).equals(Objects.requireNonNull(cs2))) { - return 0; - } - - for (int i = 0, len = Math.min(cs1.length(), cs2.length()); i < len; i++) { - final char a = cs1.charAt(i); - final char b = cs2.charAt(i); - if (a != b) { - return a - b; - } - } - - return cs1.length() - cs2.length(); - }); - break; - case BLOB: - case OBJECT: - case BOOLEAN: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE, - Type.TEXT, - Type.DATE, - Type.TIMESTAMP, - Type.STRING); - } + protected void constructPQ() { + TypeServices.BOTTOM_K_QUEUE_CONSTRUCTOR_SERVICE + .call(org.apache.tsfile.read.common.type.Type.fromTsDataType(dataType)) + .construct(this); + } + + void initializeIntQueue() { + intPQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); + } + + void initializeLongQueue() { + longPQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); + } + + void initializeFloatQueue() { + floatPQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); + } + + void initializeDoubleQueue() { + doublePQ = new PriorityQueue<>(k, Comparator.comparing(o -> -o.right)); + } + + void initializeStringQueue() { + stringPQ = + new PriorityQueue<>( + k, + (pairA, pairB) -> { + final String cs1 = pairB.right; + final String cs2 = pairA.right; + + if (Objects.requireNonNull(cs1).equals(Objects.requireNonNull(cs2))) { + return 0; + } + + for (int i = 0, len = Math.min(cs1.length(), cs2.length()); i < len; i++) { + final char a = cs1.charAt(i); + final char b = cs2.charAt(i); + if (a != b) { + return a - b; + } + } + + return cs1.length() - cs2.length(); + }); } @Override diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFChangePoints.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFChangePoints.java index dd19304384b..f573e0e320b 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFChangePoints.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFChangePoints.java @@ -19,6 +19,7 @@ package org.apache.iotdb.commons.udf.builtin; +import org.apache.iotdb.commons.udf.utils.UDFDataTypeTransformer; import org.apache.iotdb.udf.api.UDTF; import org.apache.iotdb.udf.api.access.Row; import org.apache.iotdb.udf.api.collector.PointCollector; @@ -28,19 +29,21 @@ import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.customizer.strategy.RowByRowAccessStrategy; import org.apache.iotdb.udf.api.type.Type; +import java.io.IOException; + /** * Return a series that the consecutive identical values in input series are removed (keeping only * the first one). */ public class UDTFChangePoints implements UDTF { private boolean isFirst = true; - private Type dataType; private boolean cacheBoolean; private int cacheInt; private long cacheLong; private float cacheFloat; private double cacheDouble; private String cacheString; + private TypeServices.ChangePointProcessor changePointProcessor; @Override public void validate(UDFParameterValidator validator) throws Exception { @@ -50,96 +53,99 @@ public class UDTFChangePoints implements UDTF { @Override public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) throws Exception { - dataType = parameters.getDataType(0); + Type dataType = parameters.getDataType(0); + changePointProcessor = + TypeServices.CHANGE_POINT_PROCESSOR_SERVICE.call( + UDFDataTypeTransformer.transformUDFDataTypeToReadType(dataType)); configurations.setAccessStrategy(new RowByRowAccessStrategy()).setOutputDataType(dataType); } @Override public void transform(Row row, PointCollector collector) throws Exception { - switch (dataType) { - case BOOLEAN: - if (isFirst) { - isFirst = false; - cacheBoolean = row.getBoolean(0); - collector.putBoolean(row.getTime(), cacheBoolean); - } else { - boolean rowData = row.getBoolean(0); - if (rowData != cacheBoolean) { - cacheBoolean = rowData; - collector.putBoolean(row.getTime(), cacheBoolean); - } - } - break; - case INT32: - if (isFirst) { - isFirst = false; - cacheInt = row.getInt(0); - collector.putInt(row.getTime(), cacheInt); - } else { - int rowData = row.getInt(0); - if (rowData != cacheInt) { - cacheInt = rowData; - collector.putInt(row.getTime(), cacheInt); - } - } - break; - case INT64: - if (isFirst) { - isFirst = false; - cacheLong = row.getLong(0); - collector.putLong(row.getTime(), cacheLong); - } else { - long rowData = row.getLong(0); - if (rowData != cacheLong) { - cacheLong = rowData; - collector.putLong(row.getTime(), cacheLong); - } - } - break; - case FLOAT: - if (isFirst) { - isFirst = false; - cacheFloat = row.getFloat(0); - collector.putFloat(row.getTime(), cacheFloat); - } else { - float rowData = row.getFloat(0); - if (rowData != cacheFloat) { - cacheFloat = rowData; - collector.putFloat(row.getTime(), cacheFloat); - } - } - break; - case DOUBLE: - if (isFirst) { - isFirst = false; - cacheDouble = row.getDouble(0); - collector.putDouble(row.getTime(), cacheDouble); - } else { - double rowData = row.getDouble(0); - if (rowData != cacheDouble) { - cacheDouble = rowData; - collector.putDouble(row.getTime(), cacheDouble); - } - } - break; - case TEXT: - if (isFirst) { - isFirst = false; - cacheString = row.getString(0); - collector.putString(row.getTime(), cacheString); - } else { - String rowData = row.getString(0); - if (!rowData.equals(cacheString)) { - cacheString = rowData; - collector.putString(row.getTime(), cacheString); - } - } - case STRING: - case BLOB: - case DATE: - case TIMESTAMP: - default: - break; + changePointProcessor.transform(this, row, collector); + } + + void transformBoolean(Row row, PointCollector collector) throws IOException { + if (isFirst) { + isFirst = false; + cacheBoolean = row.getBoolean(0); + collector.putBoolean(row.getTime(), cacheBoolean); + } else { + boolean rowData = row.getBoolean(0); + if (rowData != cacheBoolean) { + cacheBoolean = rowData; + collector.putBoolean(row.getTime(), cacheBoolean); + } + } + } + + void transformInt(Row row, PointCollector collector) throws IOException { + if (isFirst) { + isFirst = false; + cacheInt = row.getInt(0); + collector.putInt(row.getTime(), cacheInt); + } else { + int rowData = row.getInt(0); + if (rowData != cacheInt) { + cacheInt = rowData; + collector.putInt(row.getTime(), cacheInt); + } + } + } + + void transformLong(Row row, PointCollector collector) throws IOException { + if (isFirst) { + isFirst = false; + cacheLong = row.getLong(0); + collector.putLong(row.getTime(), cacheLong); + } else { + long rowData = row.getLong(0); + if (rowData != cacheLong) { + cacheLong = rowData; + collector.putLong(row.getTime(), cacheLong); + } + } + } + + void transformFloat(Row row, PointCollector collector) throws IOException { + if (isFirst) { + isFirst = false; + cacheFloat = row.getFloat(0); + collector.putFloat(row.getTime(), cacheFloat); + } else { + float rowData = row.getFloat(0); + if (rowData != cacheFloat) { + cacheFloat = rowData; + collector.putFloat(row.getTime(), cacheFloat); + } + } + } + + void transformDouble(Row row, PointCollector collector) throws IOException { + if (isFirst) { + isFirst = false; + cacheDouble = row.getDouble(0); + collector.putDouble(row.getTime(), cacheDouble); + } else { + double rowData = row.getDouble(0); + if (rowData != cacheDouble) { + cacheDouble = rowData; + collector.putDouble(row.getTime(), cacheDouble); + } + } + } + + void transformString(Row row, PointCollector collector) throws IOException { + if (isFirst) { + isFirst = false; + cacheString = row.getString(0); + collector.putString(row.getTime(), cacheString); + } else { + String rowData = row.getString(0); + if (!rowData.equals(cacheString)) { + cacheString = rowData; + collector.putString(row.getTime(), cacheString); + } } } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java index f3436859f2a..09ebb443a8d 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java @@ -34,6 +34,7 @@ import org.apache.iotdb.udf.api.exception.UDFParameterNotValidException; import org.apache.tsfile.block.column.Column; import org.apache.tsfile.block.column.ColumnBuilder; import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.read.common.type.Type; import org.apache.tsfile.utils.Binary; import org.apache.tsfile.utils.BytesUtils; import org.apache.tsfile.utils.DateUtils; @@ -60,14 +61,15 @@ public class UDTFConst implements UDTF { VALID_TYPES.add(TSDataType.OBJECT.name()); } - private TSDataType dataType; - private int intValue; private long longValue; private float floatValue; private double doubleValue; private boolean booleanValue; private Binary binaryValue; + private TypeServices.ConstantRowCollector rowCollector; + private TypeServices.ConstantRowMapper rowMapper; + private TypeServices.ConstantColumnValueWriter columnValueWriter; @Override public void validate(UDFParameterValidator validator) throws UDFParameterNotValidException { @@ -82,38 +84,12 @@ public class UDTFConst implements UDTF { @Override public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) { - dataType = TSDataType.valueOf(parameters.getString("type")); - switch (dataType) { - case INT32: - intValue = Integer.parseInt(parameters.getString("value")); - break; - case DATE: - intValue = DateUtils.parseDateExpressionToInt(parameters.getString("value")); - break; - case INT64: - case TIMESTAMP: - longValue = Long.parseLong(parameters.getString("value")); - break; - case FLOAT: - floatValue = Float.parseFloat(parameters.getString("value")); - break; - case DOUBLE: - doubleValue = Double.parseDouble(parameters.getString("value")); - break; - case BOOLEAN: - booleanValue = Boolean.parseBoolean(parameters.getString("value")); - break; - case TEXT: - case STRING: - binaryValue = BytesUtils.valueOf(parameters.getString("value")); - break; - case BLOB: - case OBJECT: - binaryValue = new Binary(BlobUtils.parseBlobString(parameters.getString("value"))); - break; - default: - throw new UnsupportedOperationException(); - } + TSDataType dataType = TSDataType.valueOf(parameters.getString("type")); + Type type = Type.fromTsDataType(dataType); + TypeServices.CONSTANT_PARSER_SERVICE.call(type).parse(this, parameters); + rowCollector = TypeServices.CONSTANT_ROW_COLLECTOR_SERVICE.call(type); + rowMapper = TypeServices.CONSTANT_ROW_MAPPER_SERVICE.call(type); + columnValueWriter = TypeServices.CONSTANT_COLUMN_VALUE_WRITER_SERVICE.call(type); configurations .setAccessStrategy(new MappableRowByRowAccessStrategy()) @@ -122,162 +98,93 @@ public class UDTFConst implements UDTF { @Override public void transform(Row row, PointCollector collector) throws Exception { - switch (dataType) { - case INT32: - case DATE: - collector.putInt(row.getTime(), intValue); - break; - case INT64: - case TIMESTAMP: - collector.putLong(row.getTime(), longValue); - break; - case FLOAT: - collector.putFloat(row.getTime(), floatValue); - break; - case DOUBLE: - collector.putDouble(row.getTime(), doubleValue); - break; - case BOOLEAN: - collector.putBoolean(row.getTime(), booleanValue); - break; - case TEXT: - case STRING: - case BLOB: - case OBJECT: - collector.putBinary(row.getTime(), UDFBinaryTransformer.transformToUDFBinary(binaryValue)); - break; - default: - throw new UnsupportedOperationException(); - } + rowCollector.collect(this, row, collector); } @Override public Object transform(Row row) throws IOException { - switch (dataType) { - case INT32: - case DATE: - return intValue; - case INT64: - case TIMESTAMP: - return longValue; - case FLOAT: - return floatValue; - case DOUBLE: - return doubleValue; - case BOOLEAN: - return booleanValue; - case TEXT: - case STRING: - case BLOB: - case OBJECT: - return UDFBinaryTransformer.transformToUDFBinary(binaryValue); - default: - throw new UnsupportedOperationException(); - } + return rowMapper.map(this); } @Override public void transform(Column[] columns, ColumnBuilder builder) throws Exception { - int count = columns[0].getPositionCount(); + writeConstant(columns, builder); + } - switch (dataType) { - case INT32: - case DATE: - for (int i = 0; i < count; i++) { - boolean hasWritten = false; - for (int j = 0; j < columns.length - 1; j++) { - if (!columns[j].isNull(i)) { - builder.writeInt(intValue); - hasWritten = true; - break; - } - } - if (!hasWritten) { - builder.appendNull(); - } - } - return; - case INT64: - case TIMESTAMP: - for (int i = 0; i < count; i++) { - boolean hasWritten = false; - for (int j = 0; j < columns.length - 1; j++) { - if (!columns[j].isNull(i)) { - builder.writeLong(longValue); - hasWritten = true; - break; - } - } - if (!hasWritten) { - builder.appendNull(); - } - } - return; - case FLOAT: - for (int i = 0; i < count; i++) { - boolean hasWritten = false; - for (int j = 0; j < columns.length - 1; j++) { - if (!columns[j].isNull(i)) { - builder.writeFloat(floatValue); - hasWritten = true; - break; - } - } - if (!hasWritten) { - builder.appendNull(); - } - } - return; - case DOUBLE: - for (int i = 0; i < count; i++) { - boolean hasWritten = false; - for (int j = 0; j < columns.length - 1; j++) { - if (!columns[j].isNull(i)) { - builder.writeDouble(doubleValue); - hasWritten = true; - break; - } - } - if (!hasWritten) { - builder.appendNull(); - } - } - return; - case BOOLEAN: - for (int i = 0; i < count; i++) { - boolean hasWritten = false; - for (int j = 0; j < columns.length - 1; j++) { - if (!columns[j].isNull(i)) { - builder.writeBoolean(booleanValue); - hasWritten = true; - break; - } - } - if (!hasWritten) { - builder.appendNull(); - } - } - return; - case TEXT: - case STRING: - case BLOB: - case OBJECT: - for (int i = 0; i < count; i++) { - boolean hasWritten = false; - for (int j = 0; j < columns.length - 1; j++) { - if (!columns[j].isNull(i)) { - builder.writeBinary(binaryValue); - hasWritten = true; - break; - } - } - if (!hasWritten) { - builder.appendNull(); - } + void parseInt(UDFParameters parameters) { + intValue = Integer.parseInt(parameters.getString("value")); + } + + void parseDate(UDFParameters parameters) { + intValue = DateUtils.parseDateExpressionToInt(parameters.getString("value")); + } + + void parseLong(UDFParameters parameters) { + longValue = Long.parseLong(parameters.getString("value")); + } + + void parseFloat(UDFParameters parameters) { + floatValue = Float.parseFloat(parameters.getString("value")); + } + + void parseDouble(UDFParameters parameters) { + doubleValue = Double.parseDouble(parameters.getString("value")); + } + + void parseBoolean(UDFParameters parameters) { + booleanValue = Boolean.parseBoolean(parameters.getString("value")); + } + + void parseText(UDFParameters parameters) { + binaryValue = BytesUtils.valueOf(parameters.getString("value")); + } + + void parseBlob(UDFParameters parameters) { + binaryValue = new Binary(BlobUtils.parseBlobString(parameters.getString("value"))); + } + + int intValue() { + return intValue; + } + + long longValue() { + return longValue; + } + + float floatValue() { + return floatValue; + } + + double doubleValue() { + return doubleValue; + } + + boolean booleanValue() { + return booleanValue; + } + + org.apache.iotdb.udf.api.type.Binary binaryValue() { + return UDFBinaryTransformer.transformToUDFBinary(binaryValue); + } + + Binary tsFileBinaryValue() { + return binaryValue; + } + + private void writeConstant(Column[] columns, ColumnBuilder builder) { + int count = columns[0].getPositionCount(); + for (int i = 0; i < count; i++) { + boolean hasWritten = false; + for (int j = 0; j < columns.length - 1; j++) { + if (!columns[j].isNull(i)) { + columnValueWriter.write(this, builder); + hasWritten = true; + break; } - return; - default: - throw new UnsupportedOperationException(); + } + if (!hasWritten) { + builder.appendNull(); + } } } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFContinuouslySatisfy.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFContinuouslySatisfy.java index de6232ac6db..657b75e6027 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFContinuouslySatisfy.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFContinuouslySatisfy.java @@ -32,7 +32,6 @@ import org.apache.iotdb.udf.api.exception.UDFException; import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; import org.apache.iotdb.udf.api.type.Type; -import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.utils.Pair; import java.io.IOException; @@ -40,11 +39,12 @@ import java.io.IOException; public abstract class UDTFContinuouslySatisfy implements UDTF { protected long min; protected long max; - protected TSDataType dataType; protected long satisfyValueCount; protected long satisfyValueLastTime; protected long satisfyValueStartTime; protected Pair<Long, Long> interval; + private TypeServices.ContinuouslySatisfyRowTransformer rowTransformer; + private TypeServices.ContinuouslySatisfyTerminator terminator; @Override public void validate(UDFParameterValidator validator) throws UDFException { @@ -74,7 +74,12 @@ public abstract class UDTFContinuouslySatisfy implements UDTF { satisfyValueStartTime = 0L; satisfyValueLastTime = -1L; - dataType = UDFDataTypeTransformer.transformToTsDataType(parameters.getDataType(0)); + rowTransformer = + TypeServices.CONTINUOUSLY_SATISFY_ROW_TRANSFORMER_SERVICE.call( + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0))); + terminator = + TypeServices.CONTINUOUSLY_SATISFY_TERMINATOR_SERVICE.call( + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0))); min = parameters.getLongOrDefault("min", getDefaultMin()); max = parameters.getLongOrDefault("max", getDefaultMax()); configurations.setAccessStrategy(new RowByRowAccessStrategy()).setOutputDataType(Type.INT64); @@ -83,40 +88,7 @@ public abstract class UDTFContinuouslySatisfy implements UDTF { @Override public void transform(Row row, PointCollector collector) throws IOException, UDFInputSeriesDataTypeNotValidException { - boolean needAddNewRecord; - switch (dataType) { - case INT32: - needAddNewRecord = transformInt(row.getTime(), row.getInt(0)); - break; - case INT64: - needAddNewRecord = transformLong(row.getTime(), row.getLong(0)); - break; - case FLOAT: - needAddNewRecord = transformFloat(row.getTime(), row.getFloat(0)); - break; - case DOUBLE: - needAddNewRecord = transformDouble(row.getTime(), row.getDouble(0)); - break; - case BOOLEAN: - needAddNewRecord = transformBoolean(row.getTime(), row.getBoolean(0)); - break; - case TEXT: - case STRING: - case BLOB: - case OBJECT: - case TIMESTAMP: - case DATE: - default: - // This will not happen - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } - if (needAddNewRecord) { + if (rowTransformer.transform(this, row)) { collector.putLong(interval.left, interval.right); } } @@ -215,33 +187,12 @@ public abstract class UDTFContinuouslySatisfy implements UDTF { @Override public void terminate(PointCollector collector) throws UDFInputSeriesDataTypeNotValidException, IOException { - switch (dataType) { - case INT32: - case INT64: - case FLOAT: - case DOUBLE: - case BOOLEAN: - if (satisfyValueCount > 0) { - if (getRecord() >= min && getRecord() <= max) { - collector.putLong(satisfyValueStartTime, getRecord()); - } - } - break; - case TIMESTAMP: - case DATE: - case STRING: - case BLOB: - case OBJECT: - case TEXT: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); + terminator.terminate(this, collector); + } + + void terminateSupportedType(PointCollector collector) throws IOException { + if (satisfyValueCount > 0 && getRecord() >= min && getRecord() <= max) { + collector.putLong(satisfyValueStartTime, getRecord()); } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketAggSample.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketAggSample.java index 21506698605..a98f831dd14 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketAggSample.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketAggSample.java @@ -28,9 +28,7 @@ import org.apache.iotdb.udf.api.customizer.parameter.UDFParameterValidator; import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.customizer.strategy.SlidingSizeWindowAccessStrategy; import org.apache.iotdb.udf.api.exception.UDFException; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; import org.apache.iotdb.udf.api.exception.UDFParameterNotValidException; -import org.apache.iotdb.udf.api.type.Type; import org.apache.tsfile.enums.TSDataType; @@ -40,6 +38,7 @@ public class UDTFEqualSizeBucketAggSample extends UDTFEqualSizeBucketSample { private String aggMethodType; private Aggregator aggregator; + private TypeServices.NumericWindowTransformer<UDTFEqualSizeBucketAggSample> windowTransformer; private interface Aggregator { @@ -463,40 +462,30 @@ public class UDTFEqualSizeBucketAggSample extends UDTFEqualSizeBucketSample { throw new UDFParameterNotValidException( "Illegal aggregation method. Aggregation type should be avg, min, max, sum, extreme, variance."); } + windowTransformer = + TypeServices.BUCKET_AGG_WINDOW_TRANSFORMER_SERVICE.call( + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0))); } @Override public void transform(RowWindow rowWindow, PointCollector collector) throws IOException, UDFParameterNotValidException { - switch (dataType) { - case INT32: - aggregator.aggregateInt(rowWindow, collector); - break; - case INT64: - aggregator.aggregateLong(rowWindow, collector); - break; - case FLOAT: - aggregator.aggregateFloat(rowWindow, collector); - break; - case DOUBLE: - aggregator.aggregateDouble(rowWindow, collector); - break; - case BLOB: - case OBJECT: - case TEXT: - case DATE: - case STRING: - case TIMESTAMP: - case BOOLEAN: - default: - // This will not happen - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + windowTransformer.transform(this, rowWindow, collector); + } + + void aggregateInt(RowWindow rowWindow, PointCollector collector) throws IOException { + aggregator.aggregateInt(rowWindow, collector); + } + + void aggregateLong(RowWindow rowWindow, PointCollector collector) throws IOException { + aggregator.aggregateLong(rowWindow, collector); + } + + void aggregateFloat(RowWindow rowWindow, PointCollector collector) throws IOException { + aggregator.aggregateFloat(rowWindow, collector); + } + + void aggregateDouble(RowWindow rowWindow, PointCollector collector) throws IOException { + aggregator.aggregateDouble(rowWindow, collector); } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketM4Sample.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketM4Sample.java index fea5cb694fa..b2e41eacbc7 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketM4Sample.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketM4Sample.java @@ -27,16 +27,20 @@ import org.apache.iotdb.udf.api.customizer.config.UDTFConfigurations; import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.customizer.strategy.SlidingSizeWindowAccessStrategy; import org.apache.iotdb.udf.api.exception.UDFException; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; -import org.apache.iotdb.udf.api.type.Type; + +import org.apache.tsfile.read.common.type.Type; import java.io.IOException; public class UDTFEqualSizeBucketM4Sample extends UDTFEqualSizeBucketSample { + private TypeServices.NumericWindowTransformer<UDTFEqualSizeBucketM4Sample> windowTransformer; + @Override public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) { bucketSize *= 4; + windowTransformer = + TypeServices.BUCKET_M4_WINDOW_TRANSFORMER_SERVICE.call(Type.fromTsDataType(dataType)); configurations .setAccessStrategy(new SlidingSizeWindowAccessStrategy(bucketSize)) .setOutputDataType(UDFDataTypeTransformer.transformToUDFDataType(dataType)); @@ -45,36 +49,7 @@ public class UDTFEqualSizeBucketM4Sample extends UDTFEqualSizeBucketSample { @Override public void transform(RowWindow rowWindow, PointCollector collector) throws UDFException, IOException { - switch (dataType) { - case INT32: - transformInt(rowWindow, collector); - break; - case INT64: - transformLong(rowWindow, collector); - break; - case FLOAT: - transformFloat(rowWindow, collector); - break; - case DOUBLE: - transformDouble(rowWindow, collector); - break; - case TIMESTAMP: - case BOOLEAN: - case DATE: - case STRING: - case TEXT: - case BLOB: - case OBJECT: - default: - // This will not happen - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + windowTransformer.transform(this, rowWindow, collector); } public void transformInt(RowWindow rowWindow, PointCollector collector) throws IOException { diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketOutlierSample.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketOutlierSample.java index 85490ac534e..e52d96e1fe3 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketOutlierSample.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketOutlierSample.java @@ -29,9 +29,7 @@ import org.apache.iotdb.udf.api.customizer.parameter.UDFParameterValidator; import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.customizer.strategy.SlidingSizeWindowAccessStrategy; import org.apache.iotdb.udf.api.exception.UDFException; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; import org.apache.iotdb.udf.api.exception.UDFParameterNotValidException; -import org.apache.iotdb.udf.api.type.Type; import org.apache.tsfile.utils.Pair; @@ -45,6 +43,7 @@ public class UDTFEqualSizeBucketOutlierSample extends UDTFEqualSizeBucketSample private String type; private int number; private OutlierSampler outlierSampler; + private TypeServices.NumericWindowTransformer<UDTFEqualSizeBucketOutlierSample> windowTransformer; private interface OutlierSampler { @@ -638,41 +637,31 @@ public class UDTFEqualSizeBucketOutlierSample extends UDTFEqualSizeBucketSample throw new UDFParameterNotValidException( "Illegal outlier method. Outlier type should be avg, stendis, cos or prenextdis."); } + windowTransformer = + TypeServices.BUCKET_OUTLIER_WINDOW_TRANSFORMER_SERVICE.call( + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0))); } @Override public void transform(RowWindow rowWindow, PointCollector collector) throws IOException, UDFParameterNotValidException { - switch (dataType) { - case INT32: - outlierSampler.outlierSampleInt(rowWindow, collector); - break; - case INT64: - outlierSampler.outlierSampleLong(rowWindow, collector); - break; - case FLOAT: - outlierSampler.outlierSampleFloat(rowWindow, collector); - break; - case DOUBLE: - outlierSampler.outlierSampleDouble(rowWindow, collector); - break; - case TEXT: - case BLOB: - case OBJECT: - case DATE: - case STRING: - case BOOLEAN: - case TIMESTAMP: - default: - // This will not happen - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + windowTransformer.transform(this, rowWindow, collector); + } + + void outlierSampleInt(RowWindow rowWindow, PointCollector collector) throws IOException { + outlierSampler.outlierSampleInt(rowWindow, collector); + } + + void outlierSampleLong(RowWindow rowWindow, PointCollector collector) throws IOException { + outlierSampler.outlierSampleLong(rowWindow, collector); + } + + void outlierSampleFloat(RowWindow rowWindow, PointCollector collector) throws IOException { + outlierSampler.outlierSampleFloat(rowWindow, collector); + } + + void outlierSampleDouble(RowWindow rowWindow, PointCollector collector) throws IOException { + outlierSampler.outlierSampleDouble(rowWindow, collector); } public void addToMinHeap(PriorityQueue<Pair<Integer, Double>> pq, int i, double value) { diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketRandomSample.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketRandomSample.java index 566e02e7b58..7d2ad94de08 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketRandomSample.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFEqualSizeBucketRandomSample.java @@ -26,8 +26,8 @@ import org.apache.iotdb.udf.api.collector.PointCollector; import org.apache.iotdb.udf.api.customizer.config.UDTFConfigurations; import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.customizer.strategy.SlidingSizeWindowAccessStrategy; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; -import org.apache.iotdb.udf.api.type.Type; + +import org.apache.tsfile.read.common.type.Type; import java.io.IOException; import java.security.SecureRandom; @@ -35,48 +35,20 @@ import java.security.SecureRandom; public class UDTFEqualSizeBucketRandomSample extends UDTFEqualSizeBucketSample { private SecureRandom random; + private TypeServices.NumericRowCollector rowCollector; @Override public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) { random = new SecureRandom(); + rowCollector = TypeServices.NUMERIC_ROW_COLLECTOR_SERVICE.call(Type.fromTsDataType(dataType)); configurations .setAccessStrategy(new SlidingSizeWindowAccessStrategy(bucketSize)) .setOutputDataType(UDFDataTypeTransformer.transformToUDFDataType(dataType)); } @Override - public void transform(RowWindow rowWindow, PointCollector collector) - throws IOException, UDFInputSeriesDataTypeNotValidException { + public void transform(RowWindow rowWindow, PointCollector collector) throws IOException { Row row = rowWindow.getRow(random.nextInt(rowWindow.windowSize())); - switch (dataType) { - case INT32: - collector.putInt(row.getTime(), row.getInt(0)); - break; - case INT64: - collector.putLong(row.getTime(), row.getLong(0)); - break; - case FLOAT: - collector.putFloat(row.getTime(), row.getFloat(0)); - break; - case DOUBLE: - collector.putDouble(row.getTime(), row.getDouble(0)); - break; - case BOOLEAN: - case TIMESTAMP: - case DATE: - case STRING: - case BLOB: - case OBJECT: - case TEXT: - default: - // This will not happen - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + rowCollector.collect(row, collector); } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFM4.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFM4.java index 3de46df72a0..eb7970ac0ef 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFM4.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFM4.java @@ -30,7 +30,6 @@ import org.apache.iotdb.udf.api.customizer.parameter.UDFParameters; import org.apache.iotdb.udf.api.customizer.strategy.SlidingSizeWindowAccessStrategy; import org.apache.iotdb.udf.api.customizer.strategy.SlidingTimeWindowAccessStrategy; import org.apache.iotdb.udf.api.exception.UDFException; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; import org.apache.iotdb.udf.api.exception.UDFParameterNotValidException; import org.apache.iotdb.udf.api.type.Type; @@ -60,6 +59,7 @@ public class UDTFM4 implements UDTF { protected AccessStrategy accessStrategy; protected TSDataType dataType; + private TypeServices.NumericWindowTransformer<UDTFM4> windowTransformer; public static final String WINDOW_SIZE_KEY = "windowSize"; public static final String TIME_INTERVAL_KEY = "timeInterval"; @@ -100,6 +100,9 @@ public class UDTFM4 implements UDTF { @Override public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) throws MetadataException { + windowTransformer = + TypeServices.M4_WINDOW_TRANSFORMER_SERVICE.call( + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0))); // set data type configurations.setOutputDataType(UDFDataTypeTransformer.transformToUDFDataType(dataType)); @@ -124,36 +127,7 @@ public class UDTFM4 implements UDTF { @Override public void transform(RowWindow rowWindow, PointCollector collector) throws UDFException, IOException { - switch (dataType) { - case INT32: - transformInt(rowWindow, collector); - break; - case INT64: - transformLong(rowWindow, collector); - break; - case FLOAT: - transformFloat(rowWindow, collector); - break; - case DOUBLE: - transformDouble(rowWindow, collector); - break; - case BLOB: - case OBJECT: - case DATE: - case STRING: - case TIMESTAMP: - case BOOLEAN: - case TEXT: - default: - // This will not happen - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + windowTransformer.transform(this, rowWindow, collector); } public void transformInt(RowWindow rowWindow, PointCollector collector) throws IOException { diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFSelectK.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFSelectK.java index a888c998d8d..d4c4e2bb6de 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFSelectK.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFSelectK.java @@ -49,6 +49,8 @@ public abstract class UDTFSelectK implements UDTF { protected PriorityQueue<Pair<Long, Float>> floatPQ; protected PriorityQueue<Pair<Long, Double>> doublePQ; protected PriorityQueue<Pair<Long, String>> stringPQ; + private TypeServices.SelectKRowTransformer rowTransformer; + private TypeServices.SelectKTerminator terminator; @Override public void validate(UDFParameterValidator validator) throws UDFException { @@ -77,6 +79,10 @@ public abstract class UDTFSelectK implements UDTF { k = parameters.getInt("k"); dataType = UDFDataTypeTransformer.transformToTsDataType(parameters.getDataType(0)); constructPQ(); + org.apache.tsfile.read.common.type.Type type = + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0)); + rowTransformer = TypeServices.SELECT_K_ROW_TRANSFORMER_SERVICE.call(type); + terminator = TypeServices.SELECT_K_TERMINATOR_SERVICE.call(type); configurations .setAccessStrategy(new RowByRowAccessStrategy()) .setOutputDataType(UDFDataTypeTransformer.transformToUDFDataType(dataType)); @@ -87,42 +93,7 @@ public abstract class UDTFSelectK implements UDTF { @Override public void transform(Row row, PointCollector collector) throws UDFInputSeriesDataTypeNotValidException, IOException { - switch (dataType) { - case INT32: - case DATE: - transformInt(row.getTime(), row.getInt(0)); - break; - case INT64: - case TIMESTAMP: - transformLong(row.getTime(), row.getLong(0)); - break; - case FLOAT: - transformFloat(row.getTime(), row.getFloat(0)); - break; - case DOUBLE: - transformDouble(row.getTime(), row.getDouble(0)); - break; - case TEXT: - case STRING: - transformString(row.getTime(), row.getString(0)); - break; - case BLOB: - case OBJECT: - case BOOLEAN: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE, - Type.TEXT, - Type.DATE, - Type.TIMESTAMP, - Type.STRING); - } + rowTransformer.transform(this, row); } protected abstract void transformInt(long time, int value); @@ -138,64 +109,41 @@ public abstract class UDTFSelectK implements UDTF { @Override public void terminate(PointCollector collector) throws UDFInputSeriesDataTypeNotValidException, IOException { - switch (dataType) { - case INT32: - case DATE: - for (Pair<Long, Integer> pair : - intPQ.stream().sorted(Comparator.comparing(p -> p.left)).collect(Collectors.toList())) { - collector.putInt(pair.left, pair.right); - } - break; - case INT64: - case TIMESTAMP: - for (Pair<Long, Long> pair : - longPQ.stream() - .sorted(Comparator.comparing(p -> p.left)) - .collect(Collectors.toList())) { - collector.putLong(pair.left, pair.right); - } - break; - case FLOAT: - for (Pair<Long, Float> pair : - floatPQ.stream() - .sorted(Comparator.comparing(p -> p.left)) - .collect(Collectors.toList())) { - collector.putFloat(pair.left, pair.right); - } - break; - case DOUBLE: - for (Pair<Long, Double> pair : - doublePQ.stream() - .sorted(Comparator.comparing(p -> p.left)) - .collect(Collectors.toList())) { - collector.putDouble(pair.left, pair.right); - } - break; - case TEXT: - case STRING: - for (Pair<Long, String> pair : - stringPQ.stream() - .sorted(Comparator.comparing(p -> p.left)) - .collect(Collectors.toList())) { - collector.putString(pair.left, pair.right); - } - break; - case BLOB: - case OBJECT: - case BOOLEAN: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE, - Type.TEXT, - Type.DATE, - Type.TIMESTAMP, - Type.STRING); + terminator.terminate(this, collector); + } + + void terminateInt(PointCollector collector) throws IOException { + for (Pair<Long, Integer> pair : + intPQ.stream().sorted(Comparator.comparing(p -> p.left)).collect(Collectors.toList())) { + collector.putInt(pair.left, pair.right); + } + } + + void terminateLong(PointCollector collector) throws IOException { + for (Pair<Long, Long> pair : + longPQ.stream().sorted(Comparator.comparing(p -> p.left)).collect(Collectors.toList())) { + collector.putLong(pair.left, pair.right); + } + } + + void terminateFloat(PointCollector collector) throws IOException { + for (Pair<Long, Float> pair : + floatPQ.stream().sorted(Comparator.comparing(p -> p.left)).collect(Collectors.toList())) { + collector.putFloat(pair.left, pair.right); + } + } + + void terminateDouble(PointCollector collector) throws IOException { + for (Pair<Long, Double> pair : + doublePQ.stream().sorted(Comparator.comparing(p -> p.left)).collect(Collectors.toList())) { + collector.putDouble(pair.left, pair.right); + } + } + + void terminateString(PointCollector collector) throws IOException { + for (Pair<Long, String> pair : + stringPQ.stream().sorted(Comparator.comparing(p -> p.left)).collect(Collectors.toList())) { + collector.putString(pair.left, pair.right); } } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFTopK.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFTopK.java index b8d94308ed9..fb37c91d9f4 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFTopK.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFTopK.java @@ -19,10 +19,6 @@ package org.apache.iotdb.commons.udf.builtin; -import org.apache.iotdb.commons.udf.utils.UDFDataTypeTransformer; -import org.apache.iotdb.udf.api.exception.UDFInputSeriesDataTypeNotValidException; -import org.apache.iotdb.udf.api.type.Type; - import org.apache.tsfile.utils.Pair; import java.util.Comparator; @@ -31,43 +27,30 @@ import java.util.PriorityQueue; public class UDTFTopK extends UDTFSelectK { @Override - protected void constructPQ() throws UDFInputSeriesDataTypeNotValidException { - switch (dataType) { - case INT32: - case DATE: - intPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); - break; - case INT64: - case TIMESTAMP: - longPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); - break; - case FLOAT: - floatPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); - break; - case DOUBLE: - doublePQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); - break; - case TEXT: - case STRING: - stringPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); - break; - case BOOLEAN: - case BLOB: - case OBJECT: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE, - Type.TEXT, - Type.DATE, - Type.TIMESTAMP, - Type.STRING); - } + protected void constructPQ() { + TypeServices.TOP_K_QUEUE_CONSTRUCTOR_SERVICE + .call(org.apache.tsfile.read.common.type.Type.fromTsDataType(dataType)) + .construct(this); + } + + void initializeIntQueue() { + intPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); + } + + void initializeLongQueue() { + longPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); + } + + void initializeFloatQueue() { + floatPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); + } + + void initializeDoubleQueue() { + doublePQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); + } + + void initializeStringQueue() { + stringPQ = new PriorityQueue<>(k, Comparator.comparing(o -> o.right)); } @Override
