This is an automated email from the ASF dual-hosted git repository. snuyanzin pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit 69cdb54e953e0217c5cb0c42ac50f957be01211a Author: Ramin Gharib <[email protected]> AuthorDate: Wed Jul 29 16:37:02 2026 +0200 [FLINK-40255][table-planner] Fix compiled plan serde for VARIANT type VARIANT plans and executes correctly, but serializing a plan containing it throws, so COMPILE PLAN, StatementSet#compilePlan and TableEnvironment#compilePlanSql are unusable with VARIANT, as is plan restore. Two RelDataType to LogicalType converters exist and only one knew VARIANT. FlinkTypeFactory, used during planning, handles it. LogicalRelDataTypeConverter, used only by plan serde, had no VARIANT handling at all, so a VARIANT-typed RexNode failed with "Unsupported RelDataType: VARIANT" on compile and "Logical type 'VARIANT' cannot be converted to a RelDataType" on restore. Independently, VARIANT was missing from CompactSerializationChecker, so a VARIANT in an ExecNode output type fell [...] LogicalRelDataTypeConverter now converts VARIANT in both directions, and CompactSerializationChecker reports it as compactly serializable. The latter is what fixes serialization, because the compact path already round-trips VARIANT through LogicalTypeParser. The explicit cases added to LogicalTypeJsonSerializer and LogicalTypeJsonDeserializer keep the two switches symmetric so an externally produced plan that spells VARIANT out as an object still restores. VARIANT had zero plan serde coverage. It is now exercised standalone, nested in collection and row types, as a RexNode type, and end to end by the new calc-variant restore program, whose plan carries VARIANT as the return type of PARSE_JSON and TRY_PARSE_JSON calls, as the type of an ITEM call on a VARIANT column, and in the ExecNode output type. --- .../exec/serde/LogicalTypeJsonDeserializer.java | 3 + .../exec/serde/LogicalTypeJsonSerializer.java | 2 + .../typeutils/LogicalRelDataTypeConverter.java | 8 ++ .../plan/nodes/exec/common/CalcTestPrograms.java | 35 ++++++ .../nodes/exec/serde/DataTypeJsonSerdeTest.java | 3 + .../nodes/exec/serde/LogicalTypeJsonSerdeTest.java | 6 + .../nodes/exec/serde/RelDataTypeJsonSerdeTest.java | 5 + .../nodes/exec/serde/RexNodeJsonSerdeTest.java | 6 + .../plan/nodes/exec/stream/CalcRestoreTest.java | 3 +- .../typeutils/LogicalRelDataTypeConverterTest.java | 5 + .../calc-variant/plan/calc-variant.json | 126 +++++++++++++++++++++ .../calc-variant/savepoint/_metadata | Bin 0 -> 7984 bytes 12 files changed, 201 insertions(+), 1 deletion(-) diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonDeserializer.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonDeserializer.java index 819cbe34a3f..a369885090e 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonDeserializer.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonDeserializer.java @@ -49,6 +49,7 @@ import org.apache.flink.table.types.logical.TimestampKind; import org.apache.flink.table.types.logical.TimestampType; import org.apache.flink.table.types.logical.VarBinaryType; import org.apache.flink.table.types.logical.VarCharType; +import org.apache.flink.table.types.logical.VariantType; import org.apache.flink.table.types.logical.ZonedTimestampType; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonParser; @@ -164,6 +165,8 @@ final class LogicalTypeJsonDeserializer extends StdDeserializer<LogicalType> { return deserializeStructuredType(logicalTypeNode, serdeContext); case SYMBOL: return new SymbolType<>(); + case VARIANT: + return new VariantType(); case RAW: return deserializeSpecializedRaw(logicalTypeNode, serdeContext); default: diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerializer.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerializer.java index 2857c4dbaa3..f06c38ba1b8 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerializer.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerializer.java @@ -238,6 +238,7 @@ final class LogicalTypeJsonSerializer extends StdSerializer<LogicalType> { serializeCatalogObjects); break; case SYMBOL: + case VARIANT: // type root is enough break; case RAW: @@ -544,6 +545,7 @@ final class LogicalTypeJsonSerializer extends StdSerializer<LogicalType> { case NULL: case DESCRIPTOR: case BITMAP: + case VARIANT: return true; default: // fall back to generic serialization diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverter.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverter.java index 12093cfe1e5..cf78ce82068 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverter.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverter.java @@ -59,6 +59,7 @@ import org.apache.flink.table.types.logical.TimestampType; import org.apache.flink.table.types.logical.TinyIntType; import org.apache.flink.table.types.logical.VarBinaryType; import org.apache.flink.table.types.logical.VarCharType; +import org.apache.flink.table.types.logical.VariantType; import org.apache.flink.table.types.logical.YearMonthIntervalType; import org.apache.flink.table.types.logical.YearMonthIntervalType.YearMonthResolution; import org.apache.flink.table.types.logical.ZonedTimestampType; @@ -461,6 +462,11 @@ public final class LogicalRelDataTypeConverter { return relDataTypeFactory.createSqlType(SqlTypeName.COLUMN_LIST); } + @Override + public RelDataType visit(VariantType variantType) { + return relDataTypeFactory.createSqlType(SqlTypeName.VARIANT); + } + @Override public RelDataType visit(BitmapType bitmapType) { return new BitmapRelDataType(bitmapType); @@ -588,6 +594,8 @@ public final class LogicalRelDataTypeConverter { .collect(Collectors.toList())); case COLUMN_LIST: return new DescriptorType(false); + case VARIANT: + return new VariantType(false); case STRUCTURED: case OTHER: if (relDataType instanceof StructuredRelDataType) { diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/common/CalcTestPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/common/CalcTestPrograms.java index 940c90212e4..04c86763757 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/common/CalcTestPrograms.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/common/CalcTestPrograms.java @@ -29,6 +29,8 @@ 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 org.apache.flink.types.variant.Variant; +import org.apache.flink.types.variant.VariantBuilder; import java.math.BigDecimal; import java.time.Instant; @@ -37,6 +39,8 @@ import java.time.LocalDateTime; /** {@link TableTestProgram}s for testing {@link StreamExecCalc} and {@link BatchExecCalc}. */ public class CalcTestPrograms { + private static final VariantBuilder VARIANT_BUILDER = Variant.newBuilder(); + // -------------------------------------------------------------------------------------------- // With restore data // -------------------------------------------------------------------------------------------- @@ -236,6 +240,37 @@ public class CalcTestPrograms { "INSERT INTO sink_t SELECT extract(year from current_timestamp) / a FROM t") .build(); + public static final TableTestProgram CALC_VARIANT = + TableTestProgram.of("calc-variant", "validates calc node with VARIANT type") + .setupTableSource( + SourceTestStep.newBuilder("t") + .addSchema("s STRING", "v VARIANT") + .producedBeforeRestore( + Row.of( + "{\"a\":1}", + VARIANT_BUILDER + .object() + .add("k", VARIANT_BUILDER.of(1)) + .build())) + .producedAfterRestore( + Row.of( + "{\"a\":2}", + VARIANT_BUILDER + .object() + .add("k", VARIANT_BUILDER.of(2)) + .build())) + .build()) + .setupTableSink( + SinkTestStep.newBuilder("sink_t") + .addSchema( + "parsed VARIANT", "try_parsed VARIANT", "field VARIANT") + .consumedBeforeRestore("+I[{\"a\":1}, {\"a\":1}, 1]") + .consumedAfterRestore("+I[{\"a\":2}, {\"a\":2}, 2]") + .build()) + .runSql( + "INSERT INTO sink_t SELECT PARSE_JSON(s), TRY_PARSE_JSON(s), v['k'] FROM t") + .build(); + // -------------------------------------------------------------------------------------------- // Without restore data // -------------------------------------------------------------------------------------------- diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/DataTypeJsonSerdeTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/DataTypeJsonSerdeTest.java index 3c24b472bac..bee8c249624 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/DataTypeJsonSerdeTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/DataTypeJsonSerdeTest.java @@ -59,6 +59,9 @@ public class DataTypeJsonSerdeTest { DataTypes.TIMESTAMP_LTZ(3).toInternal(), DataTypes.TIMESTAMP_LTZ(9).bridgedTo(long.class), DataTypes.BITMAP(), + DataTypes.VARIANT(), + DataTypes.VARIANT().notNull(), + DataTypes.ARRAY(DataTypes.VARIANT()), DataTypes.ROW( DataTypes.TIMESTAMP_LTZ(3).toInternal(), DataTypes.TIMESTAMP_LTZ(9).bridgedTo(long.class), diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerdeTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerdeTest.java index 0742134ebd6..c838a6ca08a 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerdeTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerdeTest.java @@ -61,6 +61,7 @@ import org.apache.flink.table.types.logical.TimestampType; import org.apache.flink.table.types.logical.TinyIntType; import org.apache.flink.table.types.logical.VarBinaryType; import org.apache.flink.table.types.logical.VarCharType; +import org.apache.flink.table.types.logical.VariantType; import org.apache.flink.table.types.logical.YearMonthIntervalType; import org.apache.flink.table.types.logical.ZonedTimestampType; import org.apache.flink.table.types.utils.DataTypeFactoryMock; @@ -267,6 +268,11 @@ public class LogicalTypeJsonSerdeTest { new MultisetType(BinaryType.ofEmptyLiteral()), new MultisetType(VarBinaryType.ofEmptyLiteral()), new BitmapType(), + new VariantType(), + new ArrayType(new VariantType()), + new MultisetType(new VariantType()), + new MapType(new VarCharType(5), new VariantType()), + RowType.of(new VariantType(), new VariantType(false)), RowType.of(new BigIntType(), new IntType(false), new VarCharType(200)), RowType.of( new LogicalType[] { diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RelDataTypeJsonSerdeTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RelDataTypeJsonSerdeTest.java index 5dc11186c26..75432d91aa7 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RelDataTypeJsonSerdeTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RelDataTypeJsonSerdeTest.java @@ -167,6 +167,11 @@ public class RelDataTypeJsonSerdeTest { FACTORY.createSqlType(SqlTypeName.VARBINARY, 1000), FACTORY.createSqlType(SqlTypeName.NULL), FACTORY.createSqlType(SqlTypeName.SYMBOL), + FACTORY.createSqlType(SqlTypeName.VARIANT), + FACTORY.createArrayType(FACTORY.createSqlType(SqlTypeName.VARIANT), -1), + FACTORY.createMapType( + FACTORY.createSqlType(SqlTypeName.VARCHAR, 5), + FACTORY.createSqlType(SqlTypeName.VARIANT)), FACTORY.createMultisetType(FACTORY.createSqlType(SqlTypeName.VARCHAR), -1), FACTORY.createArrayType(FACTORY.createSqlType(SqlTypeName.VARCHAR, 16), -1), FACTORY.createArrayType( diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RexNodeJsonSerdeTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RexNodeJsonSerdeTest.java index 047bea447ef..979ecdb5dfb 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RexNodeJsonSerdeTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RexNodeJsonSerdeTest.java @@ -830,6 +830,12 @@ public class RexNodeJsonSerdeTest { rexBuilder.makeCall( FlinkSqlOperatorTable.HASH_CODE, rexBuilder.makeInputRef(FACTORY.createSqlType(SqlTypeName.INTEGER), 1)), + rexBuilder.makeInputRef(FACTORY.createSqlType(SqlTypeName.VARIANT), 1), + // $1['f'] on a VARIANT, which also makes the call itself VARIANT-typed + rexBuilder.makeCall( + FlinkSqlOperatorTable.ITEM, + rexBuilder.makeInputRef(FACTORY.createSqlType(SqlTypeName.VARIANT), 1), + rexBuilder.makeLiteral("f")), rexBuilder.makePatternFieldRef( "test", FACTORY.createSqlType(SqlTypeName.INTEGER), 0), new RexTableArgCall( diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/CalcRestoreTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/CalcRestoreTest.java index e26f2aa1864..c4edc11e8c7 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/CalcRestoreTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/CalcRestoreTest.java @@ -43,6 +43,7 @@ public class CalcRestoreTest extends RestoreTestBase { CalcTestPrograms.CALC_UDF_SIMPLE, CalcTestPrograms.CALC_UDF_COMPLEX, CalcTestPrograms.CALC_CURRENT_TIMESTAMP, - CalcTestPrograms.COALESCE); + CalcTestPrograms.COALESCE, + CalcTestPrograms.CALC_VARIANT); } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverterTest.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverterTest.java index 92f563cc5dc..fd37601a39f 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverterTest.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverterTest.java @@ -50,6 +50,7 @@ import org.apache.flink.table.types.logical.TimestampType; import org.apache.flink.table.types.logical.TinyIntType; import org.apache.flink.table.types.logical.VarBinaryType; import org.apache.flink.table.types.logical.VarCharType; +import org.apache.flink.table.types.logical.VariantType; import org.apache.flink.table.types.logical.YearMonthIntervalType; import org.apache.flink.table.types.utils.DataTypeFactoryMock; @@ -152,7 +153,11 @@ public class LogicalRelDataTypeConverterTest { new MultisetType(BinaryType.ofEmptyLiteral()), new MultisetType(VarBinaryType.ofEmptyLiteral()), new BitmapType(), + new VariantType(), + new ArrayType(new VariantType()), + new MapType(new VarCharType(5), new VariantType()), RowType.of(new BigIntType(), new IntType(false), new VarCharType(200)), + RowType.of(new VariantType(), new VariantType(false)), RowType.of( new LogicalType[] { new BigIntType(), new IntType(false), new VarCharType(200) diff --git a/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/plan/calc-variant.json b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/plan/calc-variant.json new file mode 100644 index 00000000000..f77ba2ae4d6 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/plan/calc-variant.json @@ -0,0 +1,126 @@ +{ + "flinkVersion" : "2.4", + "nodes" : [ { + "id" : 1, + "type" : "stream-exec-table-source-scan_2", + "scanTableSource" : { + "table" : { + "identifier" : "`default_catalog`.`default_database`.`t`", + "resolvedTable" : { + "schema" : { + "columns" : [ { + "name" : "s", + "dataType" : "VARCHAR(2147483647)" + }, { + "name" : "v", + "dataType" : "VARIANT" + } ] + } + } + } + }, + "outputType" : "ROW<`s` VARCHAR(2147483647), `v` VARIANT>", + "description" : "TableSourceScan(table=[[default_catalog, default_database, t]], fields=[s, v])" + }, { + "id" : 2, + "type" : "stream-exec-calc_1", + "projection" : [ { + "kind" : "CALL", + "internalName" : "$PARSE_JSON$1", + "operands" : [ { + "kind" : "INPUT_REF", + "inputIndex" : 0, + "type" : "VARCHAR(2147483647)" + } ], + "type" : "VARIANT" + }, { + "kind" : "CALL", + "internalName" : "$TRY_PARSE_JSON$1", + "operands" : [ { + "kind" : "INPUT_REF", + "inputIndex" : 0, + "type" : "VARCHAR(2147483647)" + } ], + "type" : "VARIANT" + }, { + "kind" : "CALL", + "syntax" : "SPECIAL", + "internalName" : "$ITEM$1", + "operands" : [ { + "kind" : "INPUT_REF", + "inputIndex" : 1, + "type" : "VARIANT" + }, { + "kind" : "LITERAL", + "value" : "k", + "type" : "CHAR(1) NOT NULL" + } ], + "type" : "VARIANT" + } ], + "condition" : null, + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : "ROW<`EXPR$0` VARIANT, `EXPR$1` VARIANT, `EXPR$2` VARIANT>", + "description" : "Calc(select=[PARSE_JSON(s) AS EXPR$0, TRY_PARSE_JSON(s) AS EXPR$1, ITEM(v, 'k') AS EXPR$2])" + }, { + "id" : 3, + "type" : "stream-exec-sink_2", + "configuration" : { + "table.exec.sink.keyed-shuffle" : "AUTO", + "table.exec.sink.not-null-enforcer" : "ERROR", + "table.exec.sink.rowtime-inserter" : "ENABLED", + "table.exec.sink.type-length-enforcer" : "IGNORE", + "table.exec.sink.upsert-materialize" : "AUTO" + }, + "dynamicTableSink" : { + "table" : { + "identifier" : "`default_catalog`.`default_database`.`sink_t`", + "resolvedTable" : { + "schema" : { + "columns" : [ { + "name" : "parsed", + "dataType" : "VARIANT" + }, { + "name" : "try_parsed", + "dataType" : "VARIANT" + }, { + "name" : "field", + "dataType" : "VARIANT" + } ] + } + } + } + }, + "inputChangelogMode" : [ "INSERT" ], + "upsertMaterializeStrategy" : "VALUE", + "inputProperties" : [ { + "requiredDistribution" : { + "type" : "UNKNOWN" + }, + "damBehavior" : "PIPELINED", + "priority" : 0 + } ], + "outputType" : "ROW<`EXPR$0` VARIANT, `EXPR$1` VARIANT, `EXPR$2` VARIANT>", + "description" : "Sink(table=[default_catalog.default_database.sink_t], fields=[EXPR$0, EXPR$1, EXPR$2])" + } ], + "edges" : [ { + "source" : 1, + "target" : 2, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + }, { + "source" : 2, + "target" : 3, + "shuffle" : { + "type" : "FORWARD" + }, + "shuffleMode" : "PIPELINED" + } ] +} \ No newline at end of file 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 new file mode 100644 index 00000000000..ffd25e9a34a Binary files /dev/null and b/flink-table/flink-table-planner/src/test/resources/restore-tests/stream-exec-calc_1/calc-variant/savepoint/_metadata differ
