diff --git a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java index 72a2534b64fdb..d0d6a1389ed61 100644 --- a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java +++ b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/expressions/ValueLiteralExpression.java @@ -58,6 +58,7 @@ import java.util.Locale; import java.util.Objects; import java.util.Optional; +import java.util.UUID; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -306,6 +307,8 @@ public String asSerializableString(SqlFactory sqlFactory) { .map(EncodingUtils::escapeBackticks) .map(c -> String.format("`%s`", c)) .collect(Collectors.joining())); + case UUID: + return String.format("UUID '%s'", getValueAs(UUID.class).get()); case ARRAY: case MULTISET: case MAP: diff --git a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/BitmapType.java b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/BitmapType.java index 7dff19881d862..4e1bc706ac540 100644 --- a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/BitmapType.java +++ b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/BitmapType.java @@ -20,11 +20,9 @@ import org.apache.flink.annotation.PublicEvolving; import org.apache.flink.types.bitmap.Bitmap; -import org.apache.flink.types.bitmap.RoaringBitmapData; import java.util.Collections; import java.util.List; -import java.util.Set; /** * Data type of bitmap data. @@ -39,9 +37,6 @@ public final class BitmapType extends LogicalType { private static final long serialVersionUID = 1L; - private static final Set INPUT_OUTPUT_CONVERSION = - conversionSet(Bitmap.class.getName(), RoaringBitmapData.class.getName()); - public BitmapType(boolean isNullable) { super(isNullable, LogicalTypeRoot.BITMAP); } @@ -62,12 +57,12 @@ public String asSerializableString() { @Override public boolean supportsInputConversion(Class clazz) { - return INPUT_OUTPUT_CONVERSION.contains(clazz.getName()); + return Bitmap.class.isAssignableFrom(clazz); } @Override public boolean supportsOutputConversion(Class clazz) { - return INPUT_OUTPUT_CONVERSION.contains(clazz.getName()); + return Bitmap.class.isAssignableFrom(clazz); } @Override diff --git a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java index 623a0b5f70941..2e1411f5965ea 100644 --- a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java +++ b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java @@ -23,7 +23,6 @@ import java.util.Collections; import java.util.List; -import java.util.Set; /** * Data type of semi-structured data. @@ -38,8 +37,7 @@ @PublicEvolving public final class VariantType extends LogicalType { - private static final Set INPUT_OUTPUT_CONVERSION = - conversionSet(Variant.class.getName()); + private static final long serialVersionUID = 1L; public VariantType(boolean isNullable) { super(isNullable, LogicalTypeRoot.VARIANT); @@ -61,12 +59,12 @@ public String asSerializableString() { @Override public boolean supportsInputConversion(Class clazz) { - return INPUT_OUTPUT_CONVERSION.contains(clazz.getName()); + return Variant.class.isAssignableFrom(clazz); } @Override public boolean supportsOutputConversion(Class clazz) { - return INPUT_OUTPUT_CONVERSION.contains(clazz.getName()); + return Variant.class.isAssignableFrom(clazz); } @Override diff --git a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java index 726843ce284c6..dbcaff48b8a2e 100644 --- a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java +++ b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java @@ -29,7 +29,6 @@ import org.apache.flink.types.ColumnList; import org.apache.flink.types.Row; import org.apache.flink.types.bitmap.Bitmap; -import org.apache.flink.types.bitmap.RoaringBitmapData; import org.apache.flink.types.variant.Variant; import java.math.BigDecimal; @@ -80,9 +79,6 @@ public final class ClassDataTypeConverter { java.time.Period.class, DataTypes.INTERVAL(DataTypes.YEAR(4), DataTypes.MONTH())); addDefaultDataType(ColumnList.class, DataTypes.DESCRIPTOR()); addDefaultDataType(java.util.UUID.class, DataTypes.UUID()); - addDefaultDataType(Variant.class, DataTypes.VARIANT()); - addDefaultDataType(Bitmap.class, DataTypes.BITMAP()); - addDefaultDataType(RoaringBitmapData.class, DataTypes.BITMAP()); } private static void addDefaultDataType(Class clazz, DataType rootType) { @@ -115,6 +111,14 @@ public static Optional extractDataType(Class clazz) { return Optional.of(new AtomicDataType(new SymbolType<>(), clazz)); } + if (Variant.class.isAssignableFrom(clazz)) { + return Optional.of(DataTypes.VARIANT()); + } + + if (Bitmap.class.isAssignableFrom(clazz)) { + return Optional.of(DataTypes.BITMAP()); + } + return Optional.ofNullable(defaultDataTypes.get(clazz.getName())); } diff --git a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java index b63e83e039dbe..d3ccf06d2faa4 100644 --- a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java +++ b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java @@ -25,7 +25,7 @@ import org.apache.flink.table.types.logical.BinaryType; import org.apache.flink.table.types.logical.CharType; import org.apache.flink.table.types.logical.LogicalTypeFamily; -import org.apache.flink.types.bitmap.RoaringBitmapData; +import org.apache.flink.types.bitmap.Bitmap; import org.apache.flink.types.variant.Variant; import java.math.BigDecimal; @@ -91,12 +91,6 @@ else if (value instanceof byte[]) { // don't let the class-based extraction kick in if array elements differ return convertToArrayType((Object[]) value) .map(dt -> dt.notNull().bridgedTo(value.getClass())); - } else if (value instanceof Variant) { - // BinaryVariant is internal, so the conversion class is the Variant interface rather - // than the runtime class of the value. - return Optional.of(DataTypes.VARIANT().notNull()); - } else if (value instanceof RoaringBitmapData) { - convertedDataType = DataTypes.BITMAP(); } final Optional resultType; @@ -108,7 +102,16 @@ else if (value instanceof byte[]) { // DATE, TIME with java.sql.Time, and arrays of primitive types resultType = ClassDataTypeConverter.extractDataType(value.getClass()); } - return resultType.map(dt -> dt.notNull().bridgedTo(value.getClass())); + return resultType.map( + dt -> { + final DataType notNullDataType = dt.notNull(); + // Because they are interfaces, and we want to avoid bridgeTo internal + // conversion classes. + if (value instanceof Variant || value instanceof Bitmap) { + return notNullDataType; + } + return notNullDataType.bridgedTo(value.getClass()); + }); } private static DataType convertToCharType(String string) { diff --git a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java index fca9d0a8f668d..63748ccdc2a05 100644 --- a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java +++ b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/expressions/ExpressionTest.java @@ -47,6 +47,7 @@ import java.time.temporal.ChronoField; import java.util.LinkedHashMap; import java.util.Map; +import java.util.UUID; import java.util.stream.Stream; import static java.util.Arrays.asList; @@ -244,6 +245,16 @@ void testInstantValueLiteralExtraction() { .isEqualTo(instant.minusMillis(100)); } + @Test + void testUuidValueLiteralExtraction() { + final UUID uuid = UUID.fromString("550e8400-e29b-41d4-a716-446655440000"); + assertThat( + new ValueLiteralExpression(uuid) + .getValueAs(UUID.class) + .orElseThrow(AssertionError::new)) + .isEqualTo(uuid); + } + @Test void testOffsetDateTimeValueLiteralExtraction() { final OffsetDateTime offsetDateTime = diff --git a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java index 00f07e9375a32..0cbe9e99ac9a9 100644 --- a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java +++ b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java @@ -24,7 +24,6 @@ import org.apache.flink.table.types.utils.ClassDataTypeConverter; import org.apache.flink.types.Row; import org.apache.flink.types.bitmap.Bitmap; -import org.apache.flink.types.bitmap.RoaringBitmapData; import org.apache.flink.types.variant.Variant; import org.junit.jupiter.params.ParameterizedTest; @@ -96,8 +95,7 @@ private static Stream testData() { of(Row.class, null), of(java.util.UUID.class, DataTypes.UUID()), of(Variant.class, DataTypes.VARIANT()), - of(Bitmap.class, DataTypes.BITMAP().bridgedTo(Bitmap.class)), - of(RoaringBitmapData.class, DataTypes.BITMAP().bridgedTo(RoaringBitmapData.class))); + of(Bitmap.class, DataTypes.BITMAP())); } @ParameterizedTest(name = "[{index}] class: {0} type: {1}") diff --git a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/DataTypesTest.java b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/DataTypesTest.java index 7151f9b8c224d..eab2da22c7ac1 100644 --- a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/DataTypesTest.java +++ b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/DataTypesTest.java @@ -307,6 +307,9 @@ private static Stream testData() { .expectUnresolvedString("['EnumTypeInfo']") .lookupReturns(dummyRaw(DayOfWeek.class)) .expectResolvedDataType(dummyRaw(DayOfWeek.class)), + TestSpec.forUnresolvedDataType(DataTypes.of(UUID.class)) + .expectUnresolvedString("['java.util.UUID']") + .expectResolvedDataType(UUID()), TestSpec.forUnresolvedDataType(DataTypes.of(Variant.class)) .expectUnresolvedString("['org.apache.flink.types.variant.Variant']") .expectResolvedDataType(VARIANT()), diff --git a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java index 58e3b5bc38a6a..42b457613d6a7 100644 --- a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java +++ b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java @@ -25,7 +25,6 @@ import org.apache.flink.table.types.logical.SymbolType; import org.apache.flink.table.types.utils.ValueDataTypeConverter; import org.apache.flink.types.bitmap.Bitmap; -import org.apache.flink.types.bitmap.RoaringBitmapData; import org.apache.flink.types.variant.Variant; import org.junit.jupiter.params.ParameterizedTest; @@ -43,6 +42,7 @@ import java.time.ZoneId; import java.time.ZoneOffset; import java.time.ZonedDateTime; +import java.util.UUID; import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; @@ -121,10 +121,9 @@ private static Stream testData() { of(TimePointUnit.HOUR, new AtomicDataType(new SymbolType<>(), TimePointUnit.class)), of(new BigDecimal[0], null), of(Variant.newBuilder().of("hello"), DataTypes.VARIANT()), - of(Bitmap.empty(), DataTypes.BITMAP().bridgedTo(RoaringBitmapData.class)), - of( - Bitmap.fromArray(new int[] {1, 2}), - DataTypes.BITMAP().bridgedTo(RoaringBitmapData.class))); + of(Bitmap.empty(), DataTypes.BITMAP()), + of(Bitmap.fromArray(new int[] {1, 2}), DataTypes.BITMAP()), + of(UUID.fromString("550e8400-e29b-41d4-a716-446655440000"), DataTypes.UUID())); } @ParameterizedTest(name = "[{index}] value: {0} type: {1}") diff --git a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java index 3b7c10ab8ac63..eba215a4af94e 100644 --- a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java +++ b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java @@ -883,14 +883,6 @@ private static Stream functionSpecs() { StaticArgument.scalar("bitmap", DataTypes.BITMAP(), false)) .expectAccumulator(TypeStrategies.explicit(DataTypes.BITMAP())) .expectOutput(TypeStrategies.explicit(DataTypes.BITMAP())), - TestSpec.forScalarFunction( - "Bitmap bridged to custom Bitmap", - InvalidCustomBitmapTypeFunction1.class) - .expectErrorMessage( - "Logical type 'BITMAP' does not support a conversion from or to class 'org.apache.flink.table.types.extraction.TypeInferenceExtractorTest$CustomBitmap'."), - TestSpec.forScalarFunction("Custom Bitmap", InvalidCustomBitmapTypeFunction2.class) - .expectErrorMessage( - "Could not extract a valid type inference for function class 'org.apache.flink.table.types.extraction.TypeInferenceExtractorTest$InvalidCustomBitmapTypeFunction2'."), // --- TestSpec.forScalarFunction("Variant in scalar function", VariantTypeFunction.class) .expectStaticArgument( @@ -2752,77 +2744,4 @@ public Variant eval( return null; } } - - @FunctionHint(input = @DataTypeHint(value = "BITMAP", bridgedTo = CustomBitmap.class)) - private static class InvalidCustomBitmapTypeFunction1 extends ScalarFunction { - public Bitmap eval(Bitmap bitmap) { - return null; - } - } - - private static class InvalidCustomBitmapTypeFunction2 extends ScalarFunction { - public Bitmap eval(CustomBitmap bitmap) { - return null; - } - } - - public static class CustomBitmap implements Bitmap { - - @Override - public void add(int value) {} - - @Override - public void add(long rangeStart, long rangeEnd) {} - - @Override - public void addN(int[] values, int offset, int n) {} - - @Override - public void and(@Nullable Bitmap other) {} - - @Override - public void andNot(@Nullable Bitmap other) {} - - @Override - public void clear() {} - - @Override - public boolean contains(int value) { - return false; - } - - @Override - public int getCardinality() { - return 0; - } - - @Override - public long getLongCardinality() { - return 0; - } - - @Override - public boolean isEmpty() { - return false; - } - - @Override - public void or(@Nullable Bitmap other) {} - - @Override - public void remove(int value) {} - - @Override - public int[] toArray() { - return new int[0]; - } - - @Override - public byte[] toBytes() { - return new byte[0]; - } - - @Override - public void xor(@Nullable Bitmap other) {} - } } diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/expressions/converter/ExpressionConverter.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/expressions/converter/ExpressionConverter.java index 3a532d6f15f5d..6930c95baaa33 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/expressions/converter/ExpressionConverter.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/expressions/converter/ExpressionConverter.java @@ -74,6 +74,7 @@ import java.util.Arrays; import java.util.List; import java.util.Optional; +import java.util.UUID; import java.util.stream.Collectors; import static org.apache.flink.table.planner.typeutils.SymbolUtil.commonToCalcite; @@ -141,6 +142,13 @@ public RexNode visit(ValueLiteralExpression valueLiteral) { .collect(Collectors.toList())); } + if (type.is(LogicalTypeRoot.UUID)) { + // UUID has no generic RexBuilder#makeLiteral support, so build the literal directly. + // This also lets filter push-down round-trip a UUID predicate back into a RexNode. + return rexBuilder.makeUuidLiteral( + valueLiteral.getValueAs(UUID.class).orElseThrow(IllegalStateException::new)); + } + Object value; switch (type.getTypeRoot()) { case DECIMAL: diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GenerateUtils.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GenerateUtils.scala index d8b2490ac6135..04ea9a9c8c79e 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GenerateUtils.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/GenerateUtils.scala @@ -629,7 +629,9 @@ object GenerateUtils { s"$leftTerm.compareTo($rightTerm)" case BOOLEAN => s"($leftTerm == $rightTerm ? 0 : ($leftTerm ? 1 : -1))" - case BINARY | VARBINARY => + case BINARY | VARBINARY | UUID => + // UUID is stored as its 16-byte big-endian encoding, so it orders by the same unsigned + // byte-wise comparison as binary strings. val sortUtil = classOf[org.apache.flink.table.runtime.operators.sort.SortUtil].getCanonicalName s"$sortUtil.compareBinary($leftTerm, $rightTerm)" diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala index 4e3a79cf7a21a..8b08be95ecbb3 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/calls/ScalarOperatorGens.scala @@ -654,9 +654,10 @@ object ScalarOperatorGens { case _ => throw new CodeGenException(s"Unsupported boolean comparison '$operator'.") } } - // both sides are binary type + // both sides are binary type or UUID (both are backed by a byte[] and order by the same + // unsigned byte-wise comparison) else if ( - isBinaryString(left.resultType) && + (isBinaryString(left.resultType) || isUuid(left.resultType)) && isInteroperable(left.resultType, right.resultType) ) { val utilName = classOf[SqlFunctionUtils].getCanonicalName diff --git a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/sort/SortCodeGenerator.scala b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/sort/SortCodeGenerator.scala index f0b1dc2308423..c624e6c1d80dd 100644 --- a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/sort/SortCodeGenerator.scala +++ b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/sort/SortCodeGenerator.scala @@ -27,7 +27,7 @@ import org.apache.flink.table.planner.plan.nodes.exec.spec.SortSpec import org.apache.flink.table.runtime.generated.{GeneratedNormalizedKeyComputer, GeneratedRecordComparator, NormalizedKeyComputer, RecordComparator} import org.apache.flink.table.runtime.operators.sort.SortUtil import org.apache.flink.table.runtime.types.PlannerTypeUtils -import org.apache.flink.table.types.logical.{DecimalType, LogicalType, RowType, TimestampType} +import org.apache.flink.table.types.logical.{DecimalType, LogicalType, RowType, TimestampType, UuidType} import org.apache.flink.table.types.logical.LogicalTypeRoot._ import scala.collection.mutable @@ -425,7 +425,8 @@ class SortCodeGenerator( case DOUBLE => "Double" case BOOLEAN => "Boolean" case VARCHAR | CHAR => "String" - case VARBINARY | BINARY => "Binary" + // UUID is backed by its 16-byte big-endian encoding, so it reuses the binary accessors. + case VARBINARY | BINARY | UUID => "Binary" case DECIMAL => "Decimal" case DATE => "Int" case TIME_WITHOUT_TIME_ZONE => "Int" @@ -448,7 +449,7 @@ class SortCodeGenerator( def supportNormalizedKey(t: LogicalType): Boolean = { t.getTypeRoot match { case _ if PlannerTypeUtils.isPrimitive(t) => true - case VARCHAR | CHAR | VARBINARY | BINARY | DATE | TIME_WITHOUT_TIME_ZONE => true + case VARCHAR | CHAR | VARBINARY | BINARY | UUID | DATE | TIME_WITHOUT_TIME_ZONE => true case TIMESTAMP_WITHOUT_TIME_ZONE => // TODO: support normalize key for non-compact timestamp TimestampData.isCompact(t.asInstanceOf[TimestampType].getPrecision) @@ -474,6 +475,7 @@ class SortCodeGenerator( case DATE => 4 case TIME_WITHOUT_TIME_ZONE => 4 case DECIMAL if DecimalData.isCompact(t.asInstanceOf[DecimalType].getPrecision) => 8 + case UUID => UuidType.BYTE_LENGTH case VARCHAR | CHAR | VARBINARY | BINARY => Int.MaxValue } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/codegen/SortCodeGeneratorTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/codegen/SortCodeGeneratorTest.java index 0dbf1c098860f..335bf63262fcf 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/codegen/SortCodeGeneratorTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/codegen/SortCodeGeneratorTest.java @@ -69,6 +69,7 @@ import org.apache.flink.table.types.logical.SmallIntType; import org.apache.flink.table.types.logical.TimestampType; import org.apache.flink.table.types.logical.TinyIntType; +import org.apache.flink.table.types.logical.UuidType; import org.apache.flink.table.types.logical.VarBinaryType; import org.apache.flink.table.types.logical.VarCharType; import org.apache.flink.types.Row; @@ -109,6 +110,7 @@ class SortCodeGeneratorTest { new DecimalType(18, 2), new DecimalType(38, 18), new VarBinaryType(VarBinaryType.MAX_LENGTH), + new UuidType(), new ArrayType(new TinyIntType()), RowType.of(new IntType()), RowType.of(RowType.of(new IntType())), @@ -282,6 +284,11 @@ private Object[] generateValues(LogicalType type) { TimestampData.fromEpochMillis(rnd.nextLong(), rnd.nextInt(1000000)); } break; + case UUID: + byte[] uuidBytes = new byte[UuidType.BYTE_LENGTH]; + rnd.nextBytes(uuidBytes); + seeds[i] = uuidBytes; + break; case ARRAY: case VARBINARY: byte[] bytes = new byte[rnd.nextInt(16) + 1]; @@ -353,6 +360,9 @@ private Object value1(LogicalType type, Random rnd) { byte[] bytes2 = new byte[rnd.nextInt(7) + 1]; rnd.nextBytes(bytes2); return bytes2; + case UUID: + // all-zero: the smallest UUID under unsigned big-endian byte ordering + return new byte[UuidType.BYTE_LENGTH]; case ROW: return GenericRowData.of(new Object[] {null}); case RAW: @@ -393,6 +403,12 @@ private Object value2(LogicalType type, Random rnd) { return type instanceof VarBinaryType ? bytes : BinaryArrayData.fromPrimitiveArray(bytes); + case UUID: + // high bit set in the leading byte: greater than value1 but less than value3 only + // under unsigned ordering (signed byte ordering would rank it as the smallest) + byte[] uuidMid = new byte[UuidType.BYTE_LENGTH]; + uuidMid[0] = (byte) 0x80; + return uuidMid; case ROW: RowType rowType = (RowType) type; if (rowType.getFields().get(0).getType().getTypeRoot() == INTEGER) { @@ -440,6 +456,11 @@ private Object value3(LogicalType type, Random rnd) { return type instanceof VarBinaryType ? bytes : BinaryArrayData.fromPrimitiveArray(bytes); + case UUID: + // all-ones: the largest UUID under unsigned big-endian byte ordering + byte[] uuidMax = new byte[UuidType.BYTE_LENGTH]; + Arrays.fill(uuidMax, (byte) 0xFF); + return uuidMax; case ROW: RowType rowType = (RowType) type; if (rowType.getFields().get(0).getType().getTypeRoot() == INTEGER) { @@ -553,7 +574,8 @@ private void testInner() throws Exception { } else if (leftArray.size() > rightArray.size()) { return order ? 1 : -1; } - } else if (t.getTypeRoot() == LogicalTypeRoot.VARBINARY) { + } else if (t.getTypeRoot() == LogicalTypeRoot.VARBINARY + || t.getTypeRoot() == LogicalTypeRoot.UUID) { int comp = org.apache.flink.table.runtime.operators.sort.SortUtil .compareBinary((byte[]) first, (byte[]) second); @@ -617,7 +639,7 @@ private void testInner() throws Exception { RowData.createFieldGetter(keyTypes[j], keys[j]); Object o1 = fieldGetter.getFieldOrNull(data.get(i)); Object o2 = fieldGetter.getFieldOrNull(result.get(i)); - if (keyTypes[j] instanceof VarBinaryType) { + if (keyTypes[j] instanceof VarBinaryType || keyTypes[j] instanceof UuidType) { assertThat((byte[]) o2).as(msg).isEqualTo((byte[]) o1); } else if (keyTypes[j] instanceof TypeInformationRawType) { assertThat((RawValueData) o1) diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/CastFunctionITCase.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/CastFunctionITCase.java index 7dbbc379c67dd..09ada039cbb9b 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/CastFunctionITCase.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/functions/CastFunctionITCase.java @@ -538,9 +538,6 @@ private static List variantArrayCasts() { } private static List uuidCasts() { - // A UUID has no Table API literal, so a UUID source is produced from a string cast. The - // exhaustive value and lenient-parse matrix lives in CastRulesTest. - final String literal = "UUID '" + DEFAULT_UUID_STRING + "'"; return Arrays.asList( // A string or a binary value casts to a UUID. CastTestSpecBuilder.testCastTo(UUID()) @@ -577,20 +574,16 @@ private static List uuidCasts() { TableRuntimeException.class, "requires exactly 16 bytes") .testResult($("f1").tryCast(UUID()), "TRY_CAST(f1 AS UUID)", null, UUID()), - // A UUID casts to its canonical string and to its 16-byte encoding. - TestSetSpec.forExpression("Cast a UUID to a string and to bytes") - .onFieldsWithData("unused") - .andDataTypes(STRING()) - .testResult( - lit(DEFAULT_UUID_STRING).cast(UUID()).cast(STRING()), - "CAST(" + literal + " AS STRING)", - DEFAULT_UUID_STRING, - STRING().notNull()) - .testResult( - lit(DEFAULT_UUID_STRING).cast(UUID()).cast(BYTES()), - "CAST(" + literal + " AS BYTES)", - DEFAULT_UUID_BYTES, - BYTES().notNull())); + // A UUID casts to its canonical string. + CastTestSpecBuilder.testCastTo(STRING()) + .fromCase(UUID(), DEFAULT_UUID, DEFAULT_UUID_STRING) + .fromCase(UUID(), null, null) + .build(), + // A UUID casts to its 16-byte encoding. + CastTestSpecBuilder.testCastTo(BYTES()) + .fromCase(UUID(), DEFAULT_UUID, DEFAULT_UUID_BYTES) + .fromCase(UUID(), null, null) + .build()); } private static List variantRowCasts() { diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/UuidTestPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/UuidTestPrograms.java new file mode 100644 index 0000000000000..68770d766b77f --- /dev/null +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/UuidTestPrograms.java @@ -0,0 +1,284 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.planner.plan.nodes.exec; + +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.api.ValidationException; +import org.apache.flink.table.test.program.SinkTestStep; +import org.apache.flink.table.test.program.SourceTestStep; +import org.apache.flink.table.test.program.TableTestProgram; +import org.apache.flink.types.Row; + +import java.util.Map; +import java.util.UUID; + +/** {@link TableTestProgram}s for the {@link DataTypes#UUID()} type. */ +public class UuidTestPrograms { + + private static final String LITERAL_A = "550e8400-e29b-41d4-a716-446655440000"; + private static final String LITERAL_B = "f47ac10b-58cc-4372-a567-0e02b2c3d479"; + private static final UUID UUID_A = UUID.fromString(LITERAL_A); + private static final UUID UUID_B = UUID.fromString(LITERAL_B); + + // Boundary values around the signed/unsigned byte split at 0x7f/0x80. Under the unsigned + // big-endian ordering of the type: ZERO < HIGH_BIT < MAX. A signed byte comparison would + // instead rank HIGH_BIT (0x80 = -128) as the smallest, so these values make the two orderings + // disagree and let the tests pin down the unsigned semantics. + private static final String LITERAL_ZERO = "00000000-0000-0000-0000-000000000000"; + private static final String LITERAL_HIGH_BIT = "80000000-0000-0000-0000-000000000000"; + private static final String LITERAL_MAX = "ffffffff-ffff-ffff-ffff-ffffffffffff"; + private static final UUID UUID_ZERO = UUID.fromString(LITERAL_ZERO); + private static final UUID UUID_HIGH_BIT = UUID.fromString(LITERAL_HIGH_BIT); + private static final UUID UUID_MAX = UUID.fromString(LITERAL_MAX); + + private static SourceTestStep singleRowDriver() { + return SourceTestStep.newBuilder("t").addSchema("d INT").producedValues(Row.of(1)).build(); + } + + public static final TableTestProgram UUID_SOURCE_SINK = + TableTestProgram.of("uuid-source-sink", "round-trips a UUID column including null") + .setupTableSource( + SourceTestStep.newBuilder("t") + .addSchema("id UUID") + .producedValues(Row.of(UUID_A), Row.of(UUID_B), new Row(1)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("id UUID") + .consumedValues(Row.of(UUID_A), Row.of(UUID_B), new Row(1)) + .build()) + .runSql("INSERT INTO sink_t SELECT id FROM t") + .build(); + + public static final TableTestProgram UUID_LITERAL = + TableTestProgram.of("uuid-literal", "materializes a UUID literal") + .setupTableSource(singleRowDriver()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("u UUID") + .consumedValues(Row.of(UUID_A)) + .build()) + .runSql("INSERT INTO sink_t SELECT UUID '" + LITERAL_A + "' FROM t") + .build(); + + public static final TableTestProgram UUID_ARRAY = + TableTestProgram.of("uuid-array", "reads UUID array elements") + .setupTableSource(singleRowDriver()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("arr ARRAY") + .consumedValues(Row.of((Object) new UUID[] {UUID_A, UUID_B})) + .build()) + .runSql( + "INSERT INTO sink_t SELECT ARRAY[UUID '" + + LITERAL_A + + "', UUID '" + + LITERAL_B + + "'] FROM t") + .build(); + + public static final TableTestProgram UUID_MAP = + TableTestProgram.of("uuid-map", "reads a UUID map value") + .setupTableSource(singleRowDriver()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("m MAP") + .consumedValues(Row.of(Map.of("a", UUID_A))) + .build()) + .runSql("INSERT INTO sink_t SELECT MAP['a', UUID '" + LITERAL_A + "'] FROM t") + .build(); + + public static final TableTestProgram UUID_NESTED_ROW = + TableTestProgram.of("uuid-nested-row", "reads a UUID field of a nested row") + .setupTableSource(singleRowDriver()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("r ROW") + .consumedValues(Row.of(Row.of(UUID_A, 42))) + .build()) + .runSql("INSERT INTO sink_t SELECT (UUID '" + LITERAL_A + "', 42) FROM t") + .build(); + + private static SourceTestStep comparisonSource() { + return SourceTestStep.newBuilder("t") + .addSchema("a UUID", "b UUID") + .producedValues( + Row.of(UUID_ZERO, UUID_ZERO), + Row.of(UUID_ZERO, UUID_HIGH_BIT), + Row.of(UUID_HIGH_BIT, UUID_ZERO), + Row.of(UUID_HIGH_BIT, UUID_MAX), + Row.of(UUID_MAX, UUID_MAX)) + .build(); + } + + public static final TableTestProgram UUID_EQUALITY = + TableTestProgram.of("uuid-equality", "compares two UUID columns for equality") + .setupTableSource(comparisonSource()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("a UUID") + .consumedValues(Row.of(UUID_ZERO), Row.of(UUID_MAX)) + .build()) + .runSql("INSERT INTO sink_t SELECT a FROM t WHERE a = b") + .build(); + + // Both 0x00 < 0x80 and 0x80 < 0xff hold under unsigned ordering. A signed byte comparison would + // drop the first pair (0x00 < 0x80 becomes 0 < -128), so this pins down the unsigned semantics. + public static final TableTestProgram UUID_COMPARISON = + TableTestProgram.of( + "uuid-comparison", + "compares two UUID columns using unsigned big-endian byte ordering") + .setupTableSource(comparisonSource()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("a UUID", "b UUID") + .consumedValues( + Row.of(UUID_ZERO, UUID_HIGH_BIT), + Row.of(UUID_HIGH_BIT, UUID_MAX)) + .build()) + .runSql("INSERT INTO sink_t SELECT a, b FROM t WHERE a < b") + .build(); + + public static final TableTestProgram UUID_LITERAL_FILTER = + TableTestProgram.of( + "uuid-literal-filter", + "filters UUID rows with a range predicate against a UUID literal") + .setupTableSource( + SourceTestStep.newBuilder("t") + .addSchema("id UUID") + .producedValues( + Row.of(UUID_ZERO), + Row.of(UUID_HIGH_BIT), + Row.of(UUID_MAX), + new Row(1)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("id UUID") + .consumedValues(Row.of(UUID_HIGH_BIT), Row.of(UUID_MAX)) + .build()) + .runSql( + "INSERT INTO sink_t SELECT id FROM t " + + "WHERE id > UUID '7fffffff-ffff-ffff-ffff-ffffffffffff'") + .build(); + + public static final TableTestProgram UUID_INVALID_LITERAL = + TableTestProgram.of("uuid-invalid-literal", "rejects a malformed UUID literal") + .setupTableSource(singleRowDriver()) + .runFailingSql( + "SELECT UUID 'abcd' FROM t", + ValidationException.class, + "Invalid UUID string: abcd") + .build(); + + private static SourceTestStep groupingSource() { + return SourceTestStep.newBuilder("t") + .addSchema("id UUID") + .producedValues( + Row.of(UUID_ZERO), + Row.of(UUID_ZERO), + Row.of(UUID_A), + Row.of(UUID_MAX), + Row.of(UUID_MAX), + Row.of(UUID_MAX)) + .build(); + } + + // The row number encodes the sort position, so a set comparison still pins down the order: + // ZERO < UUID_A (0x55) < HIGH_BIT (0x80) < MAX (0xff) holds only under the unsigned ordering. A + // bounded Top-N (rn <= 4) is used instead of a plain ORDER BY so the program is also valid in + // streaming mode, where a global sort on a non-time attribute is not supported. + public static final TableTestProgram UUID_ORDER_BY = + TableTestProgram.of( + "uuid-order-by", + "orders UUID rows by the unsigned big-endian byte comparison") + .setupTableSource( + SourceTestStep.newBuilder("t") + .addSchema("id UUID") + .producedValues( + Row.of(UUID_MAX), + Row.of(UUID_ZERO), + Row.of(UUID_HIGH_BIT), + Row.of(UUID_A)) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id UUID", "rn BIGINT", "PRIMARY KEY (rn) NOT ENFORCED") + .testMaterializedData() + .consumedValues( + Row.of(UUID_ZERO, 1L), + Row.of(UUID_A, 2L), + Row.of(UUID_HIGH_BIT, 3L), + Row.of(UUID_MAX, 4L)) + .build()) + .runSql( + "INSERT INTO sink_t SELECT id, rn FROM " + + "(SELECT id, ROW_NUMBER() OVER (ORDER BY id) AS rn FROM t) " + + "WHERE rn <= 4") + .build(); + + public static final TableTestProgram UUID_GROUP_BY = + TableTestProgram.of("uuid-group-by", "groups rows by a UUID key") + .setupTableSource(groupingSource()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "id UUID", + "cnt BIGINT", + "PRIMARY KEY (id) NOT ENFORCED") + .testMaterializedData() + .consumedValues( + Row.of(UUID_ZERO, 2L), + Row.of(UUID_A, 1L), + Row.of(UUID_MAX, 3L)) + .build()) + .runSql("INSERT INTO sink_t SELECT id, COUNT(*) FROM t GROUP BY id") + .build(); + + public static final TableTestProgram UUID_JOIN = + TableTestProgram.of("uuid-join", "joins two tables on a UUID key") + .setupTableSource( + SourceTestStep.newBuilder("l") + .addSchema("id UUID") + .producedValues( + Row.of(UUID_ZERO), + Row.of(UUID_A), + Row.of(UUID_HIGH_BIT), + Row.of(UUID_MAX)) + .build()) + .setupTableSource( + SourceTestStep.newBuilder("r") + .addSchema("id UUID", "tag STRING") + .producedValues( + Row.of(UUID_A, "a"), + Row.of(UUID_HIGH_BIT, "h"), + Row.of(UUID_MAX, "m")) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema("id UUID", "tag STRING") + .consumedValues( + Row.of(UUID_A, "a"), + Row.of(UUID_HIGH_BIT, "h"), + Row.of(UUID_MAX, "m")) + .build()) + .runSql("INSERT INTO sink_t SELECT l.id, r.tag FROM l JOIN r ON l.id = r.id") + .build(); +} diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/batch/UuidSemanticTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/batch/UuidSemanticTest.java new file mode 100644 index 0000000000000..61c45df678e3a --- /dev/null +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/batch/UuidSemanticTest.java @@ -0,0 +1,43 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.table.planner.plan.nodes.exec.batch; + +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.planner.plan.nodes.exec.UuidTestPrograms; +import org.apache.flink.table.planner.plan.nodes.exec.testutils.BatchSemanticTestBase; +import org.apache.flink.table.test.program.TableTestProgram; + +import java.util.List; + +/** + * Batch semantic tests for the {@link DataTypes#UUID()} type as an ordering, grouping and join key. + */ +public class UuidSemanticTest extends BatchSemanticTestBase { + + @Override + public List programs() { + return List.of( + UuidTestPrograms.UUID_EQUALITY, + UuidTestPrograms.UUID_COMPARISON, + UuidTestPrograms.UUID_LITERAL_FILTER, + UuidTestPrograms.UUID_ORDER_BY, + UuidTestPrograms.UUID_GROUP_BY, + UuidTestPrograms.UUID_JOIN); + } +} diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/UuidSemanticTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/UuidSemanticTest.java index debd187340326..a2d2313d7f610 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/UuidSemanticTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/UuidSemanticTest.java @@ -19,6 +19,7 @@ package org.apache.flink.table.planner.plan.nodes.exec.stream; import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.planner.plan.nodes.exec.UuidTestPrograms; import org.apache.flink.table.planner.plan.nodes.exec.testutils.SemanticTestBase; import org.apache.flink.table.test.program.TableTestProgram; @@ -35,6 +36,12 @@ public List programs() { UuidTestPrograms.UUID_ARRAY, UuidTestPrograms.UUID_MAP, UuidTestPrograms.UUID_NESTED_ROW, + UuidTestPrograms.UUID_EQUALITY, + UuidTestPrograms.UUID_COMPARISON, + UuidTestPrograms.UUID_LITERAL_FILTER, + UuidTestPrograms.UUID_ORDER_BY, + UuidTestPrograms.UUID_GROUP_BY, + UuidTestPrograms.UUID_JOIN, UuidTestPrograms.UUID_INVALID_LITERAL); } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/UuidTestPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/UuidTestPrograms.java deleted file mode 100644 index 4e60704520c6e..0000000000000 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/UuidTestPrograms.java +++ /dev/null @@ -1,115 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.flink.table.planner.plan.nodes.exec.stream; - -import org.apache.flink.table.api.DataTypes; -import org.apache.flink.table.api.ValidationException; -import org.apache.flink.table.test.program.SinkTestStep; -import org.apache.flink.table.test.program.SourceTestStep; -import org.apache.flink.table.test.program.TableTestProgram; -import org.apache.flink.types.Row; - -import java.util.Map; -import java.util.UUID; - -/** {@link TableTestProgram}s for the {@link DataTypes#UUID()} type. */ -public class UuidTestPrograms { - - private static final String LITERAL_A = "550e8400-e29b-41d4-a716-446655440000"; - private static final String LITERAL_B = "f47ac10b-58cc-4372-a567-0e02b2c3d479"; - private static final UUID UUID_A = UUID.fromString(LITERAL_A); - private static final UUID UUID_B = UUID.fromString(LITERAL_B); - - private static SourceTestStep singleRowDriver() { - return SourceTestStep.newBuilder("t").addSchema("d INT").producedValues(Row.of(1)).build(); - } - - static final TableTestProgram UUID_SOURCE_SINK = - TableTestProgram.of("uuid-source-sink", "round-trips a UUID column including null") - .setupTableSource( - SourceTestStep.newBuilder("t") - .addSchema("id UUID") - .producedValues(Row.of(UUID_A), Row.of(UUID_B), new Row(1)) - .build()) - .setupTableSink( - SinkTestStep.newBuilder("sink_t") - .addSchema("id UUID") - .consumedValues(Row.of(UUID_A), Row.of(UUID_B), new Row(1)) - .build()) - .runSql("INSERT INTO sink_t SELECT id FROM t") - .build(); - - static final TableTestProgram UUID_LITERAL = - TableTestProgram.of("uuid-literal", "materializes a UUID literal") - .setupTableSource(singleRowDriver()) - .setupTableSink( - SinkTestStep.newBuilder("sink_t") - .addSchema("u UUID") - .consumedValues(Row.of(UUID_A)) - .build()) - .runSql("INSERT INTO sink_t SELECT UUID '" + LITERAL_A + "' FROM t") - .build(); - - static final TableTestProgram UUID_ARRAY = - TableTestProgram.of("uuid-array", "reads UUID array elements") - .setupTableSource(singleRowDriver()) - .setupTableSink( - SinkTestStep.newBuilder("sink_t") - .addSchema("arr ARRAY") - .consumedValues(Row.of((Object) new UUID[] {UUID_A, UUID_B})) - .build()) - .runSql( - "INSERT INTO sink_t SELECT ARRAY[UUID '" - + LITERAL_A - + "', UUID '" - + LITERAL_B - + "'] FROM t") - .build(); - - static final TableTestProgram UUID_MAP = - TableTestProgram.of("uuid-map", "reads a UUID map value") - .setupTableSource(singleRowDriver()) - .setupTableSink( - SinkTestStep.newBuilder("sink_t") - .addSchema("m MAP") - .consumedValues(Row.of(Map.of("a", UUID_A))) - .build()) - .runSql("INSERT INTO sink_t SELECT MAP['a', UUID '" + LITERAL_A + "'] FROM t") - .build(); - - static final TableTestProgram UUID_NESTED_ROW = - TableTestProgram.of("uuid-nested-row", "reads a UUID field of a nested row") - .setupTableSource(singleRowDriver()) - .setupTableSink( - SinkTestStep.newBuilder("sink_t") - .addSchema("r ROW") - .consumedValues(Row.of(Row.of(UUID_A, 42))) - .build()) - .runSql("INSERT INTO sink_t SELECT (UUID '" + LITERAL_A + "', 42) FROM t") - .build(); - - static final TableTestProgram UUID_INVALID_LITERAL = - TableTestProgram.of("uuid-invalid-literal", "rejects a malformed UUID literal") - .setupTableSource(singleRowDriver()) - .runFailingSql( - "SELECT UUID 'abcd' FROM t", - ValidationException.class, - "Invalid UUID string: abcd") - .build(); -} diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/savepoint/_metadata b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/savepoint/_metadata index ffd25e9a34a23..b66e951dcce56 100644 Binary files a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/savepoint/_metadata and b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/savepoint/_metadata differ diff --git a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/data/util/DataFormatConverters.java b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/data/util/DataFormatConverters.java index 756ca6293c867..8d3b18d2f27b8 100644 --- a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/data/util/DataFormatConverters.java +++ b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/data/util/DataFormatConverters.java @@ -68,7 +68,6 @@ import org.apache.flink.table.utils.DateTimeUtils; import org.apache.flink.types.Row; import org.apache.flink.types.bitmap.Bitmap; -import org.apache.flink.types.bitmap.RoaringBitmapData; import org.apache.flink.types.variant.Variant; import org.apache.commons.lang3.ArrayUtils; @@ -158,8 +157,7 @@ public class DataFormatConverters { DataTypes.INTERVAL(DataTypes.SECOND(3)).bridgedTo(long.class), LongConverter.INSTANCE); - t2C.put(DataTypes.BITMAP().bridgedTo(Bitmap.class), BitmapConverter.INSTANCE); - t2C.put(DataTypes.BITMAP().bridgedTo(RoaringBitmapData.class), BitmapConverter.INSTANCE); + t2C.put(DataTypes.BITMAP(), BitmapConverter.INSTANCE); TYPE_TO_CONVERTER = Collections.unmodifiableMap(t2C); } @@ -758,15 +756,6 @@ public static final class BitmapConverter extends IdentityConverter { private BitmapConverter() {} - @Override - Bitmap toInternalImpl(Bitmap value) { - if (!(value instanceof RoaringBitmapData)) { - throw new UnsupportedOperationException( - "Unsupported bitmap type: " + value.getClass().getSimpleName() + "."); - } - return value; - } - @Override Bitmap toExternalImpl(RowData row, int column) { return row.getBitmap(column); diff --git a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/typeutils/TypeCheckUtils.java b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/typeutils/TypeCheckUtils.java index c469696d247d6..814c21c2603e9 100644 --- a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/typeutils/TypeCheckUtils.java +++ b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/typeutils/TypeCheckUtils.java @@ -139,6 +139,10 @@ public static boolean isStructuredType(LogicalType type) { return type.getTypeRoot() == STRUCTURED_TYPE; } + public static boolean isUuid(LogicalType type) { + return type.getTypeRoot() == UUID; + } + private static boolean isVariantType(LogicalType type) { return type.getTypeRoot() == VARIANT; } @@ -147,10 +151,6 @@ private static boolean isBitmapType(LogicalType type) { return type.getTypeRoot() == BITMAP; } - private static boolean isUuidType(LogicalType type) { - return type.getTypeRoot() == UUID; - } - public static boolean isComparable(LogicalType type) { return !isRaw(type) && !isMap(type) @@ -159,8 +159,7 @@ public static boolean isComparable(LogicalType type) { && !isArray(type) && !isStructuredType(type) && !isVariantType(type) - && !isBitmapType(type) - && !isUuidType(type); + && !isBitmapType(type); } public static boolean isMutable(LogicalType type) { diff --git a/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/BitmapBitmapConverter.java b/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/BitmapBitmapConverter.java deleted file mode 100644 index cde40ec057eca..0000000000000 --- a/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/BitmapBitmapConverter.java +++ /dev/null @@ -1,44 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.flink.table.data.conversion; - -import org.apache.flink.annotation.Internal; -import org.apache.flink.types.bitmap.Bitmap; -import org.apache.flink.types.bitmap.RoaringBitmapData; - -/** Converter for {@link Bitmap} type that only accepts {@link RoaringBitmapData} as input. */ -@Internal -public class BitmapBitmapConverter implements DataStructureConverter { - - private static final long serialVersionUID = 1L; - - @Override - public Bitmap toInternal(Bitmap external) { - if (!(external instanceof RoaringBitmapData)) { - throw new UnsupportedOperationException( - "Unsupported bitmap type: " + external.getClass().getSimpleName() + "."); - } - return external; - } - - @Override - public Bitmap toExternal(Bitmap internal) { - return internal; - } -} diff --git a/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/DataStructureConverters.java b/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/DataStructureConverters.java index 9d7ffdaea3d9d..61459d070968d 100644 --- a/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/DataStructureConverters.java +++ b/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/DataStructureConverters.java @@ -31,9 +31,6 @@ import org.apache.flink.table.types.logical.LogicalType; import org.apache.flink.table.types.logical.LogicalTypeRoot; import org.apache.flink.types.Row; -import org.apache.flink.types.bitmap.Bitmap; -import org.apache.flink.types.bitmap.RoaringBitmapData; -import org.apache.flink.types.variant.Variant; import java.math.BigDecimal; import java.time.Duration; @@ -200,9 +197,6 @@ public final class DataStructureConverters { putConverter(LogicalTypeRoot.RAW, RawValueData.class, identity()); putConverter(LogicalTypeRoot.UUID, UUID.class, constructor(UuidUuidConverter::new)); putConverter(LogicalTypeRoot.UUID, byte[].class, identity()); - putConverter(LogicalTypeRoot.VARIANT, Variant.class, identity()); - putConverter(LogicalTypeRoot.BITMAP, Bitmap.class, constructor(BitmapBitmapConverter::new)); - putConverter(LogicalTypeRoot.BITMAP, RoaringBitmapData.class, identity()); } /** Returns a converter for the given {@link DataType}. */ @@ -242,6 +236,10 @@ public static DataStructureConverter getConverter(DataType dataT return StructuredObjectConverter.create(dataType); case RAW: return RawObjectConverter.create(dataType); + case VARIANT: + case BITMAP: + // The conversion class is already validated by supports{Input,Output}Conversion. + return IdentityConverter.INSTANCE; default: throw new TableException("Could not find converter for data type: " + dataType); } @@ -257,7 +255,7 @@ private static void putConverter( } private static DataStructureConverterFactory identity() { - return constructor(IdentityConverter::new); + return constructor(() -> IdentityConverter.INSTANCE); } private static DataStructureConverterFactory constructor( diff --git a/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/IdentityConverter.java b/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/IdentityConverter.java index 0e9cfeff9f244..b8f008cb9176b 100644 --- a/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/IdentityConverter.java +++ b/flink-table/flink-table-type-utils/src/main/java/org/apache/flink/table/data/conversion/IdentityConverter.java @@ -24,6 +24,8 @@ @Internal public class IdentityConverter implements DataStructureConverter { + public static final IdentityConverter INSTANCE = new IdentityConverter<>(); + private static final long serialVersionUID = 1L; @Override