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 58a688e8f831c1c8bfb95d16fe5f068a14d78926 Author: Tian Jiang <[email protected]> AuthorDate: Thu Aug 27 19:08:13 2026 +0800 multiple refactors --- .../iotdb/commons/udf/builtin/TypeServices.java | 51 +++++++ .../iotdb/commons/udf/builtin/UDTFInRange.java | 162 ++------------------- .../apache/iotdb/commons/udf/builtin/UDTFMath.java | 149 ++----------------- .../iotdb/commons/udf/builtin/UDTFOnOff.java | 153 ++----------------- 4 files changed, 92 insertions(+), 423 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 f663171abae..2957e953ea2 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 @@ -18,10 +18,13 @@ 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.collector.PointCollector; 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.read.common.type.service.TypeService; import java.io.IOException; @@ -175,16 +178,54 @@ final class TypeServices { }; }; + // UDF Row has no generic numeric accessor, so bind its primitive getter once during beforeStart. + static final TypeService<NumericRowReader> NUMERIC_ROW_READER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32 -> row -> row.getInt(0); + case INT64 -> row -> row.getLong(0); + case FLOAT -> row -> row.getFloat(0); + case DOUBLE -> row -> row.getDouble(0); + case BOOLEAN, DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + row -> { + throw invalidNumericDataType(type); + }; + }; + + // TsFile Type provides primitive numeric conversion for every supported column implementation. + static final TypeService<NumericColumnReader> NUMERIC_COLUMN_READER_SERVICE = + type -> + switch (type.getTypeEnum()) { + case INT32, INT64, FLOAT, DOUBLE -> type::getDouble; + case BOOLEAN, DATE, TIMESTAMP, TEXT, STRING, BLOB, OBJECT, ROW, UNKNOWN, VECTOR -> + (column, position) -> { + throw invalidNumericDataType(type); + }; + }; + static { VALUE_TREND_READER_SERVICE.check(); VALUE_DIFFERENCE_OPERATOR_SERVICE.check(); NON_NEGATIVE_VALUE_DIFFERENCE_OPERATOR_SERVICE.check(); DERIVATIVE_OPERATOR_SERVICE.check(); NON_NEGATIVE_DERIVATIVE_OPERATOR_SERVICE.check(); + NUMERIC_ROW_READER_SERVICE.check(); + NUMERIC_COLUMN_READER_SERVICE.check(); } private TypeServices() {} + private static UDFInputSeriesDataTypeNotValidException invalidNumericDataType( + org.apache.tsfile.read.common.type.Type type) { + return new UDFInputSeriesDataTypeNotValidException( + 0, + UDFDataTypeTransformer.transformReadTypeToUDFDataType(type), + Type.INT32, + Type.INT64, + Type.FLOAT, + Type.DOUBLE); + } + @FunctionalInterface interface PreviousValueReader { void read(UDTFValueTrend target, Row row) @@ -203,4 +244,14 @@ final class TypeServices { UDTFValueTrend target, long time, Row row, PointCollector collector, double timeDelta) throws UDFInputSeriesDataTypeNotValidException, IOException; } + + @FunctionalInterface + interface NumericRowReader { + double read(Row row) throws UDFInputSeriesDataTypeNotValidException, IOException; + } + + @FunctionalInterface + interface NumericColumnReader { + double read(Column column, int position) throws UDFInputSeriesDataTypeNotValidException; + } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFInRange.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFInRange.java index 4c17a270697..4b3299f538c 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFInRange.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFInRange.java @@ -29,19 +29,18 @@ 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.MappableRowByRowAccessStrategy; 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.block.column.Column; import org.apache.tsfile.block.column.ColumnBuilder; -import org.apache.tsfile.enums.TSDataType; import java.io.IOException; public class UDTFInRange implements UDTF { - protected TSDataType dataType; protected double upper; protected double lower; + private TypeServices.NumericRowReader rowReader; + private TypeServices.NumericColumnReader columnReader; @Override public void validate(UDFParameterValidator validator) throws UDFException { @@ -62,46 +61,20 @@ public class UDTFInRange implements UDTF { throws MetadataException { upper = parameters.getDouble("upper"); lower = parameters.getDouble("lower"); - dataType = UDFDataTypeTransformer.transformToTsDataType(parameters.getDataType(0)); + org.apache.tsfile.read.common.type.Type type = + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0)); + rowReader = TypeServices.NUMERIC_ROW_READER_SERVICE.call(type); + columnReader = TypeServices.NUMERIC_COLUMN_READER_SERVICE.call(type); configurations .setAccessStrategy(new MappableRowByRowAccessStrategy()) .setOutputDataType(Type.BOOLEAN); } @Override - public void transform(Row row, PointCollector collector) - throws UDFInputSeriesDataTypeNotValidException, IOException { + public void transform(Row row, PointCollector collector) throws IOException { long time = row.getTime(); - switch (dataType) { - case INT32: - collector.putBoolean(time, row.getInt(0) >= lower && upper >= row.getInt(0)); - break; - case INT64: - collector.putBoolean(time, row.getLong(0) >= lower && upper >= row.getLong(0)); - break; - case FLOAT: - collector.putBoolean(time, row.getFloat(0) >= lower && upper >= row.getFloat(0)); - break; - case DOUBLE: - collector.putBoolean(time, row.getDouble(0) >= lower && upper >= row.getDouble(0)); - 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); - } + double value = rowReader.read(row); + collector.putBoolean(time, value >= lower && upper >= value); } @Override @@ -109,124 +82,21 @@ public class UDTFInRange implements UDTF { if (row.isNull(0)) { return null; } - switch (dataType) { - case INT32: - return row.getInt(0) >= lower && upper >= row.getInt(0); - case INT64: - return row.getLong(0) >= lower && upper >= row.getLong(0); - case FLOAT: - return row.getFloat(0) >= lower && upper >= row.getFloat(0); - case DOUBLE: - return row.getDouble(0) >= lower && upper >= row.getDouble(0); - 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); - } + double value = rowReader.read(row); + return value >= lower && upper >= value; } @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 TEXT: - 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); - } - } - - private void transformInt(Column[] columns, ColumnBuilder builder) { - int[] inputs = columns[0].getInts(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - boolean res = inputs[i] >= lower && upper >= inputs[i]; - builder.writeBoolean(res); - } - } - } - - private void transformLong(Column[] columns, ColumnBuilder builder) { - long[] inputs = columns[0].getLongs(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - boolean res = inputs[i] >= lower && upper >= inputs[i]; - builder.writeBoolean(res); - } - } - } - - private void transformFloat(Column[] columns, ColumnBuilder builder) { - float[] inputs = columns[0].getFloats(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - boolean res = inputs[i] >= lower && upper >= inputs[i]; - builder.writeBoolean(res); - } - } - } - - private void transformDouble(Column[] columns, ColumnBuilder builder) { - double[] inputs = columns[0].getDoubles(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); + Column column = columns[0]; + boolean[] isNulls = column.isNull(); + int count = column.getPositionCount(); for (int i = 0; i < count; i++) { if (isNulls[i]) { builder.appendNull(); } else { - boolean res = inputs[i] >= lower && upper >= inputs[i]; - builder.writeBoolean(res); + double value = columnReader.read(column, i); + builder.writeBoolean(value >= lower && upper >= value); } } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFMath.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFMath.java index 947cbdb2c01..08179e24c54 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFMath.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFMath.java @@ -48,6 +48,8 @@ public abstract class UDTFMath implements UDTF { protected Transformer transformer; protected TSDataType dataType; + private TypeServices.NumericRowReader rowReader; + private TypeServices.NumericColumnReader columnReader; @Override public void validate(UDFParameterValidator validator) throws UDFException { @@ -60,6 +62,10 @@ public abstract class UDTFMath implements UDTF { public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) throws MetadataException { dataType = UDFDataTypeTransformer.transformToTsDataType(parameters.getDataType(0)); + org.apache.tsfile.read.common.type.Type type = + org.apache.tsfile.read.common.type.Type.fromTsDataType(dataType); + rowReader = TypeServices.NUMERIC_ROW_READER_SERVICE.call(type); + columnReader = TypeServices.NUMERIC_COLUMN_READER_SERVICE.call(type); configurations .setAccessStrategy(new MappableRowByRowAccessStrategy()) .setOutputDataType(Type.DOUBLE); @@ -71,37 +77,7 @@ public abstract class UDTFMath implements UDTF { @Override public void transform(Row row, PointCollector collector) throws UDFInputSeriesDataTypeNotValidException, IOException { - long time = row.getTime(); - switch (dataType) { - case INT32: - collector.putDouble(time, transformer.transform(row.getInt(0))); - break; - case INT64: - collector.putDouble(time, transformer.transform(row.getLong(0))); - break; - case FLOAT: - collector.putDouble(time, transformer.transform(row.getFloat(0))); - break; - case DOUBLE: - collector.putDouble(time, transformer.transform(row.getDouble(0))); - break; - case BOOLEAN: - case TEXT: - case STRING: - case TIMESTAMP: - case DATE: - case BLOB: - case OBJECT: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + collector.putDouble(row.getTime(), transformer.transform(rowReader.read(row))); } @Override @@ -109,120 +85,19 @@ public abstract class UDTFMath implements UDTF { if (row.isNull(0)) { return null; } - switch (dataType) { - case INT32: - return transformer.transform(row.getInt(0)); - case INT64: - return transformer.transform(row.getLong(0)); - case FLOAT: - return transformer.transform(row.getFloat(0)); - case DOUBLE: - return transformer.transform(row.getDouble(0)); - case DATE: - case BLOB: - case OBJECT: - case STRING: - case TIMESTAMP: - case TEXT: - case BOOLEAN: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + return transformer.transform(rowReader.read(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 TEXT: - case BOOLEAN: - case STRING: - case TIMESTAMP: - case BLOB: - case OBJECT: - case DATE: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } - } - - private void transformInt(Column[] columns, ColumnBuilder builder) { - int[] inputs = columns[0].getInts(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - builder.writeDouble(transformer.transform(inputs[i])); - } - } - } - - private void transformLong(Column[] columns, ColumnBuilder builder) { - long[] inputs = columns[0].getLongs(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - builder.writeDouble(transformer.transform(inputs[i])); - } - } - } - - private void transformFloat(Column[] columns, ColumnBuilder builder) { - float[] inputs = columns[0].getFloats(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - builder.writeDouble(transformer.transform(inputs[i])); - } - } - } - - private void transformDouble(Column[] columns, ColumnBuilder builder) { - double[] inputs = columns[0].getDoubles(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); + Column column = columns[0]; + boolean[] isNulls = column.isNull(); + int count = column.getPositionCount(); for (int i = 0; i < count; i++) { if (isNulls[i]) { builder.appendNull(); } else { - builder.writeDouble(transformer.transform(inputs[i])); + builder.writeDouble(transformer.transform(columnReader.read(column, i))); } } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFOnOff.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFOnOff.java index b0c95bd0d23..2095d07172e 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFOnOff.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFOnOff.java @@ -29,12 +29,10 @@ 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.MappableRowByRowAccessStrategy; 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.block.column.Column; import org.apache.tsfile.block.column.ColumnBuilder; -import org.apache.tsfile.enums.TSDataType; import java.io.IOException; @@ -42,7 +40,8 @@ public class UDTFOnOff implements UDTF { protected double threshold; - private TSDataType dataType; + private TypeServices.NumericRowReader rowReader; + private TypeServices.NumericColumnReader columnReader; @Override public void validate(UDFParameterValidator validator) throws UDFException { @@ -56,45 +55,19 @@ public class UDTFOnOff implements UDTF { public void beforeStart(UDFParameters parameters, UDTFConfigurations configurations) throws MetadataException { threshold = parameters.getDouble("threshold"); - dataType = UDFDataTypeTransformer.transformToTsDataType(parameters.getDataType(0)); + org.apache.tsfile.read.common.type.Type type = + UDFDataTypeTransformer.transformUDFDataTypeToReadType(parameters.getDataType(0)); + rowReader = TypeServices.NUMERIC_ROW_READER_SERVICE.call(type); + columnReader = TypeServices.NUMERIC_COLUMN_READER_SERVICE.call(type); configurations .setAccessStrategy(new MappableRowByRowAccessStrategy()) .setOutputDataType(Type.BOOLEAN); } @Override - public void transform(Row row, PointCollector collector) - throws UDFInputSeriesDataTypeNotValidException, IOException { + public void transform(Row row, PointCollector collector) throws IOException { long time = row.getTime(); - switch (dataType) { - case INT32: - collector.putBoolean(time, row.getInt(0) >= threshold); - break; - case INT64: - collector.putBoolean(time, (row.getLong(0) >= threshold)); - break; - case FLOAT: - collector.putBoolean(time, (row.getFloat(0) >= threshold)); - break; - case DOUBLE: - collector.putBoolean(time, (row.getDouble(0) >= threshold)); - break; - case DATE: - case BLOB: - case STRING: - case TIMESTAMP: - case TEXT: - case BOOLEAN: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + collector.putBoolean(time, rowReader.read(row) >= threshold); } @Override @@ -102,119 +75,19 @@ public class UDTFOnOff implements UDTF { if (row.isNull(0)) { return null; } - switch (dataType) { - case INT32: - return row.getInt(0) >= threshold; - case INT64: - return row.getLong(0) >= threshold; - case FLOAT: - return row.getFloat(0) >= threshold; - case DOUBLE: - return row.getDouble(0) >= threshold; - case TEXT: - case BOOLEAN: - case STRING: - case TIMESTAMP: - case BLOB: - case DATE: - default: - // This will not happen. - throw new UDFInputSeriesDataTypeNotValidException( - 0, - UDFDataTypeTransformer.transformToUDFDataType(dataType), - Type.INT32, - Type.INT64, - Type.FLOAT, - Type.DOUBLE); - } + return rowReader.read(row) >= threshold; } @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 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); - } - } - - private void transformInt(Column[] columns, ColumnBuilder builder) { - int[] inputs = columns[0].getInts(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - builder.writeBoolean(inputs[i] >= threshold); - } - } - } - - private void transformLong(Column[] columns, ColumnBuilder builder) { - long[] inputs = columns[0].getLongs(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - builder.writeBoolean(inputs[i] >= threshold); - } - } - } - - private void transformFloat(Column[] columns, ColumnBuilder builder) { - float[] inputs = columns[0].getFloats(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); - for (int i = 0; i < count; i++) { - if (isNulls[i]) { - builder.appendNull(); - } else { - builder.writeBoolean(inputs[i] >= threshold); - } - } - } - - private void transformDouble(Column[] columns, ColumnBuilder builder) { - double[] inputs = columns[0].getDoubles(); - boolean[] isNulls = columns[0].isNull(); - - int count = columns[0].getPositionCount(); + Column column = columns[0]; + boolean[] isNulls = column.isNull(); + int count = column.getPositionCount(); for (int i = 0; i < count; i++) { if (isNulls[i]) { builder.appendNull(); } else { - builder.writeBoolean(inputs[i] >= threshold); + builder.writeBoolean(columnReader.read(column, i) >= threshold); } } }
