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 652cb02d0144233c2e87944645bbe88678b289ee Author: Tian Jiang <[email protected]> AuthorDate: Wed Aug 26 14:32:14 2026 +0800 multiple refactors --- .../relational/MergeSortFullOuterJoinOperator.java | 6 +-- .../calc/plan/planner/TableOperatorGenerator.java | 46 ++-------------------- .../org/apache/iotdb/calc/utils/TypeServices.java | 19 +++++++++ .../apache/iotdb/calc/utils/TypeServicesTest.java | 20 ++++++++++ 4 files changed, 46 insertions(+), 45 deletions(-) diff --git a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/source/relational/MergeSortFullOuterJoinOperator.java b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/source/relational/MergeSortFullOuterJoinOperator.java index 278d87ea021..cfb2c587a40 100644 --- a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/source/relational/MergeSortFullOuterJoinOperator.java +++ b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/execution/operator/source/relational/MergeSortFullOuterJoinOperator.java @@ -23,6 +23,7 @@ import org.apache.iotdb.calc.execution.operator.CommonOperatorContext; import org.apache.iotdb.calc.execution.operator.Operator; import org.apache.iotdb.calc.execution.operator.process.join.merge.comparator.JoinKeyComparator; import org.apache.iotdb.calc.plan.planner.CommonOperatorUtils; +import org.apache.iotdb.calc.utils.TypeServices; import org.apache.iotdb.commons.queryengine.execution.MemoryEstimationHelper; import org.apache.tsfile.block.column.Column; @@ -32,7 +33,6 @@ import org.apache.tsfile.read.common.block.TsBlock; import org.apache.tsfile.utils.RamUsageEstimator; import java.util.List; -import java.util.function.BiFunction; public class MergeSortFullOuterJoinOperator extends AbstractMergeSortJoinOperator { private static final long INSTANCE_SIZE = @@ -41,7 +41,7 @@ public class MergeSortFullOuterJoinOperator extends AbstractMergeSortJoinOperato // stores last row matched join criteria, only used in outer join private TsBlock lastMatchedRightBlock = null; private final int[] lastMatchedBlockPositions; - private final List<BiFunction<Column, Integer, Column>> updateLastMatchedRowFunctions; + private final List<TypeServices.ColumnRowFunction> updateLastMatchedRowFunctions; public MergeSortFullOuterJoinOperator( CommonOperatorContext operatorContext, @@ -53,7 +53,7 @@ public class MergeSortFullOuterJoinOperator extends AbstractMergeSortJoinOperato int[] rightOutputSymbolIdx, List<JoinKeyComparator> joinKeyComparators, List<TSDataType> dataTypes, - List<BiFunction<Column, Integer, Column>> updateLastMatchedRowFunctions) { + List<TypeServices.ColumnRowFunction> updateLastMatchedRowFunctions) { super( operatorContext, leftChild, diff --git a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/TableOperatorGenerator.java b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/TableOperatorGenerator.java index 86786404bc4..8b66d4eef27 100644 --- a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/TableOperatorGenerator.java +++ b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/TableOperatorGenerator.java @@ -108,6 +108,7 @@ import org.apache.iotdb.calc.plan.relational.planner.CastToStringLiteralVisitor; import org.apache.iotdb.calc.plan.relational.planner.CastToTimestampLiteralVisitor; import org.apache.iotdb.calc.transformation.dag.column.ColumnTransformer; import org.apache.iotdb.calc.transformation.dag.column.leaf.LeafColumnTransformer; +import org.apache.iotdb.calc.utils.TypeServices; import org.apache.iotdb.calc.utils.datastructure.SortKey; import org.apache.iotdb.common.rpc.thrift.TAggregationType; import org.apache.iotdb.commons.exception.SemanticException; @@ -183,12 +184,6 @@ import org.apache.tsfile.common.conf.TSFileConfig; import org.apache.tsfile.common.conf.TSFileDescriptor; import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.read.common.block.TsBlock; -import org.apache.tsfile.read.common.block.column.BinaryColumn; -import org.apache.tsfile.read.common.block.column.BooleanColumn; -import org.apache.tsfile.read.common.block.column.DoubleColumn; -import org.apache.tsfile.read.common.block.column.FloatColumn; -import org.apache.tsfile.read.common.block.column.IntColumn; -import org.apache.tsfile.read.common.block.column.LongColumn; import org.apache.tsfile.read.common.block.column.RunLengthEncodedColumn; import org.apache.tsfile.read.common.type.Type; import org.apache.tsfile.utils.Binary; @@ -205,7 +200,6 @@ import java.util.Map; import java.util.Optional; import java.util.OptionalInt; import java.util.Set; -import java.util.function.BiFunction; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -1175,7 +1169,9 @@ public abstract class TableOperatorGenerator< rightOutputSymbolIdx, JoinKeyComparatorFactory.getComparators(joinKeyTypes, true), dataTypes, - joinKeyTypes.stream().map(this::buildUpdateLastRowFunction).collect(Collectors.toList())); + joinKeyTypes.stream() + .map(TypeServices.UPDATE_LAST_ROW_SERVICE::call) + .collect(Collectors.toList())); } else if (requireNonNull(node.getJoinType()) == JoinNode.JoinType.LEFT) { CommonOperatorContext operatorContext = addOperatorContext( @@ -1210,40 +1206,6 @@ public abstract class TableOperatorGenerator< } } - protected BiFunction<Column, Integer, Column> buildUpdateLastRowFunction(Type joinKeyType) { - switch (joinKeyType.getTypeEnum()) { - case INT32: - return (inputColumn, rowIndex) -> - new IntColumn( - 1, Optional.empty(), new int[] {inputColumn.getInt(rowIndex)}, TSDataType.INT32); - case DATE: - return (inputColumn, rowIndex) -> - new IntColumn( - 1, Optional.empty(), new int[] {inputColumn.getInt(rowIndex)}, TSDataType.DATE); - case INT64: - case TIMESTAMP: - return (inputColumn, rowIndex) -> - new LongColumn(1, Optional.empty(), new long[] {inputColumn.getLong(rowIndex)}); - case FLOAT: - return (inputColumn, rowIndex) -> - new FloatColumn(1, Optional.empty(), new float[] {inputColumn.getFloat(rowIndex)}); - case DOUBLE: - return (inputColumn, rowIndex) -> - new DoubleColumn(1, Optional.empty(), new double[] {inputColumn.getDouble(rowIndex)}); - case BOOLEAN: - return (inputColumn, rowIndex) -> - new BooleanColumn( - 1, Optional.empty(), new boolean[] {inputColumn.getBoolean(rowIndex)}); - case STRING: - case TEXT: - case BLOB: - return (inputColumn, rowIndex) -> - new BinaryColumn(1, Optional.empty(), new Binary[] {inputColumn.getBinary(rowIndex)}); - default: - throw new UnsupportedOperationException(CalcMessages.UNSUPPORTED_DATA_TYPE + joinKeyType); - } - } - @Override public Operator visitSemiJoin(SemiJoinNode node, C context) { List<TSDataType> dataTypes = getOutputColumnTypes(node, context.getTableTypeProvider()); diff --git a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/utils/TypeServices.java b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/utils/TypeServices.java index 5980317b17c..19583424db4 100644 --- a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/utils/TypeServices.java +++ b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/utils/TypeServices.java @@ -339,6 +339,19 @@ public class TypeServices { }; }; + /** + * Creates a single-row column containing the value at the requested row. Delegating the write to + * {@link Type} preserves type-specific column metadata (for example, DATE) without another + * TsDataType dispatch. + */ + public static final TypeService<ColumnRowFunction> UPDATE_LAST_ROW_SERVICE = + type -> + (column, rowIndex) -> { + ColumnBuilder columnBuilder = type.createColumnBuilder(1); + type.write(columnBuilder, column, rowIndex); + return columnBuilder.build(); + }; + public static final TypeService<Function<TsPrimitiveType, Object>> PRIMITIVE_TYPE_VALUE_EXTRACTOR_SERVICE = type -> @@ -579,6 +592,7 @@ public class TypeServices { JOIN_KEY_COMPARATOR_SERVICE.check(); GREATEST_COLUMN_TRANSFORMER_SERVICE.check(); LEAST_COLUMN_TRANSFORMER_SERVICE.check(); + UPDATE_LAST_ROW_SERVICE.check(); MERGE_SORT_COMPARATOR_SERVICE.check(); MEMORY_USAGE_OF_ONE_MERGE_SORT_KEY_SERVICE.check(); MEMORY_USAGE_OF_ONE_SERIALIZABLE_ROW_FIELD_SERVICE.check(); @@ -618,6 +632,11 @@ public class TypeServices { TSEncoding getDefaultTextEncoding(); } + @FunctionalInterface + public interface ColumnRowFunction { + Column apply(Column column, int rowIndex); + } + @FunctionalInterface public interface ColumnToDoubleConverter { double convert(Column column, int position); diff --git a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/utils/TypeServicesTest.java b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/utils/TypeServicesTest.java index 39ca2a1a2b4..33d3723550c 100644 --- a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/utils/TypeServicesTest.java +++ b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/utils/TypeServicesTest.java @@ -151,6 +151,26 @@ public class TypeServicesTest { .apply(columnTransformers)); } + @Test + public void testUpdateLastRowServiceUsesTypeColumnWriter() { + Column input = new IntColumn(2, Optional.empty(), new int[] {10, 20}, TSDataType.DATE); + Column output = + TypeServices.UPDATE_LAST_ROW_SERVICE + .call(Type.fromTsDataType(TSDataType.DATE)) + .apply(input, 1); + + assertEquals(1, output.getPositionCount()); + assertEquals(20, output.getInt(0)); + assertEquals(TSDataType.DATE, output.getDataType()); + + Assert.assertThrows( + RuntimeException.class, + () -> + TypeServices.UPDATE_LAST_ROW_SERVICE + .call(Type.fromTsDataType(TSDataType.VECTOR)) + .apply(input, 0)); + } + // Covers every supported RANGE-frame type and guards native integer overflow and long precision. @Test public void testRangeFrameComparatorPreservesNativeArithmetic() {
