Skip to content

Commit 69cdb54

Browse files
raminqafsnuyanzin
authored andcommitted
[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 back to generic serialization and failed with "Unable to serialize logical type 'VARIANT'". 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.
1 parent 115c4fb commit 69cdb54

12 files changed

Lines changed: 201 additions & 1 deletion

File tree

flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonDeserializer.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@
4949
import org.apache.flink.table.types.logical.TimestampType;
5050
import org.apache.flink.table.types.logical.VarBinaryType;
5151
import org.apache.flink.table.types.logical.VarCharType;
52+
import org.apache.flink.table.types.logical.VariantType;
5253
import org.apache.flink.table.types.logical.ZonedTimestampType;
5354

5455
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonParser;
@@ -164,6 +165,8 @@ private static LogicalType deserializeFromRoot(
164165
return deserializeStructuredType(logicalTypeNode, serdeContext);
165166
case SYMBOL:
166167
return new SymbolType<>();
168+
case VARIANT:
169+
return new VariantType();
167170
case RAW:
168171
return deserializeSpecializedRaw(logicalTypeNode, serdeContext);
169172
default:

flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerializer.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -238,6 +238,7 @@ private static void serializeTypeWithGenericSerialization(
238238
serializeCatalogObjects);
239239
break;
240240
case SYMBOL:
241+
case VARIANT:
241242
// type root is enough
242243
break;
243244
case RAW:
@@ -544,6 +545,7 @@ protected Boolean defaultMethod(LogicalType logicalType) {
544545
case NULL:
545546
case DESCRIPTOR:
546547
case BITMAP:
548+
case VARIANT:
547549
return true;
548550
default:
549551
// fall back to generic serialization

flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverter.java

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,7 @@
5959
import org.apache.flink.table.types.logical.TinyIntType;
6060
import org.apache.flink.table.types.logical.VarBinaryType;
6161
import org.apache.flink.table.types.logical.VarCharType;
62+
import org.apache.flink.table.types.logical.VariantType;
6263
import org.apache.flink.table.types.logical.YearMonthIntervalType;
6364
import org.apache.flink.table.types.logical.YearMonthIntervalType.YearMonthResolution;
6465
import org.apache.flink.table.types.logical.ZonedTimestampType;
@@ -461,6 +462,11 @@ public RelDataType visit(DescriptorType descriptorType) {
461462
return relDataTypeFactory.createSqlType(SqlTypeName.COLUMN_LIST);
462463
}
463464

465+
@Override
466+
public RelDataType visit(VariantType variantType) {
467+
return relDataTypeFactory.createSqlType(SqlTypeName.VARIANT);
468+
}
469+
464470
@Override
465471
public RelDataType visit(BitmapType bitmapType) {
466472
return new BitmapRelDataType(bitmapType);
@@ -588,6 +594,8 @@ private static LogicalType toLogicalTypeNotNull(
588594
.collect(Collectors.toList()));
589595
case COLUMN_LIST:
590596
return new DescriptorType(false);
597+
case VARIANT:
598+
return new VariantType(false);
591599
case STRUCTURED:
592600
case OTHER:
593601
if (relDataType instanceof StructuredRelDataType) {

flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/common/CalcTestPrograms.java

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,8 @@
2929
import org.apache.flink.table.test.program.SourceTestStep;
3030
import org.apache.flink.table.test.program.TableTestProgram;
3131
import org.apache.flink.types.Row;
32+
import org.apache.flink.types.variant.Variant;
33+
import org.apache.flink.types.variant.VariantBuilder;
3234

3335
import java.math.BigDecimal;
3436
import java.time.Instant;
@@ -37,6 +39,8 @@
3739
/** {@link TableTestProgram}s for testing {@link StreamExecCalc} and {@link BatchExecCalc}. */
3840
public class CalcTestPrograms {
3941

42+
private static final VariantBuilder VARIANT_BUILDER = Variant.newBuilder();
43+
4044
// --------------------------------------------------------------------------------------------
4145
// With restore data
4246
// --------------------------------------------------------------------------------------------
@@ -236,6 +240,37 @@ public class CalcTestPrograms {
236240
"INSERT INTO sink_t SELECT extract(year from current_timestamp) / a FROM t")
237241
.build();
238242

243+
public static final TableTestProgram CALC_VARIANT =
244+
TableTestProgram.of("calc-variant", "validates calc node with VARIANT type")
245+
.setupTableSource(
246+
SourceTestStep.newBuilder("t")
247+
.addSchema("s STRING", "v VARIANT")
248+
.producedBeforeRestore(
249+
Row.of(
250+
"{\"a\":1}",
251+
VARIANT_BUILDER
252+
.object()
253+
.add("k", VARIANT_BUILDER.of(1))
254+
.build()))
255+
.producedAfterRestore(
256+
Row.of(
257+
"{\"a\":2}",
258+
VARIANT_BUILDER
259+
.object()
260+
.add("k", VARIANT_BUILDER.of(2))
261+
.build()))
262+
.build())
263+
.setupTableSink(
264+
SinkTestStep.newBuilder("sink_t")
265+
.addSchema(
266+
"parsed VARIANT", "try_parsed VARIANT", "field VARIANT")
267+
.consumedBeforeRestore("+I[{\"a\":1}, {\"a\":1}, 1]")
268+
.consumedAfterRestore("+I[{\"a\":2}, {\"a\":2}, 2]")
269+
.build())
270+
.runSql(
271+
"INSERT INTO sink_t SELECT PARSE_JSON(s), TRY_PARSE_JSON(s), v['k'] FROM t")
272+
.build();
273+
239274
// --------------------------------------------------------------------------------------------
240275
// Without restore data
241276
// --------------------------------------------------------------------------------------------

flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/DataTypeJsonSerdeTest.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,9 @@ private static Stream<DataType> testDataTypeSerde() {
5959
DataTypes.TIMESTAMP_LTZ(3).toInternal(),
6060
DataTypes.TIMESTAMP_LTZ(9).bridgedTo(long.class),
6161
DataTypes.BITMAP(),
62+
DataTypes.VARIANT(),
63+
DataTypes.VARIANT().notNull(),
64+
DataTypes.ARRAY(DataTypes.VARIANT()),
6265
DataTypes.ROW(
6366
DataTypes.TIMESTAMP_LTZ(3).toInternal(),
6467
DataTypes.TIMESTAMP_LTZ(9).bridgedTo(long.class),

flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/LogicalTypeJsonSerdeTest.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@
6161
import org.apache.flink.table.types.logical.TinyIntType;
6262
import org.apache.flink.table.types.logical.VarBinaryType;
6363
import org.apache.flink.table.types.logical.VarCharType;
64+
import org.apache.flink.table.types.logical.VariantType;
6465
import org.apache.flink.table.types.logical.YearMonthIntervalType;
6566
import org.apache.flink.table.types.logical.ZonedTimestampType;
6667
import org.apache.flink.table.types.utils.DataTypeFactoryMock;
@@ -267,6 +268,11 @@ private static List<LogicalType> testLogicalTypeSerde() {
267268
new MultisetType(BinaryType.ofEmptyLiteral()),
268269
new MultisetType(VarBinaryType.ofEmptyLiteral()),
269270
new BitmapType(),
271+
new VariantType(),
272+
new ArrayType(new VariantType()),
273+
new MultisetType(new VariantType()),
274+
new MapType(new VarCharType(5), new VariantType()),
275+
RowType.of(new VariantType(), new VariantType(false)),
270276
RowType.of(new BigIntType(), new IntType(false), new VarCharType(200)),
271277
RowType.of(
272278
new LogicalType[] {

flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RelDataTypeJsonSerdeTest.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -167,6 +167,11 @@ public static List<RelDataType> testRelDataTypeSerde() {
167167
FACTORY.createSqlType(SqlTypeName.VARBINARY, 1000),
168168
FACTORY.createSqlType(SqlTypeName.NULL),
169169
FACTORY.createSqlType(SqlTypeName.SYMBOL),
170+
FACTORY.createSqlType(SqlTypeName.VARIANT),
171+
FACTORY.createArrayType(FACTORY.createSqlType(SqlTypeName.VARIANT), -1),
172+
FACTORY.createMapType(
173+
FACTORY.createSqlType(SqlTypeName.VARCHAR, 5),
174+
FACTORY.createSqlType(SqlTypeName.VARIANT)),
170175
FACTORY.createMultisetType(FACTORY.createSqlType(SqlTypeName.VARCHAR), -1),
171176
FACTORY.createArrayType(FACTORY.createSqlType(SqlTypeName.VARCHAR, 16), -1),
172177
FACTORY.createArrayType(

flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/serde/RexNodeJsonSerdeTest.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -830,6 +830,12 @@ private static Stream<RexNode> testRexNodeSerde() {
830830
rexBuilder.makeCall(
831831
FlinkSqlOperatorTable.HASH_CODE,
832832
rexBuilder.makeInputRef(FACTORY.createSqlType(SqlTypeName.INTEGER), 1)),
833+
rexBuilder.makeInputRef(FACTORY.createSqlType(SqlTypeName.VARIANT), 1),
834+
// $1['f'] on a VARIANT, which also makes the call itself VARIANT-typed
835+
rexBuilder.makeCall(
836+
FlinkSqlOperatorTable.ITEM,
837+
rexBuilder.makeInputRef(FACTORY.createSqlType(SqlTypeName.VARIANT), 1),
838+
rexBuilder.makeLiteral("f")),
833839
rexBuilder.makePatternFieldRef(
834840
"test", FACTORY.createSqlType(SqlTypeName.INTEGER), 0),
835841
new RexTableArgCall(

flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/CalcRestoreTest.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ public List<TableTestProgram> programs() {
4343
CalcTestPrograms.CALC_UDF_SIMPLE,
4444
CalcTestPrograms.CALC_UDF_COMPLEX,
4545
CalcTestPrograms.CALC_CURRENT_TIMESTAMP,
46-
CalcTestPrograms.COALESCE);
46+
CalcTestPrograms.COALESCE,
47+
CalcTestPrograms.CALC_VARIANT);
4748
}
4849
}

flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/typeutils/LogicalRelDataTypeConverterTest.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,7 @@
5050
import org.apache.flink.table.types.logical.TinyIntType;
5151
import org.apache.flink.table.types.logical.VarBinaryType;
5252
import org.apache.flink.table.types.logical.VarCharType;
53+
import org.apache.flink.table.types.logical.VariantType;
5354
import org.apache.flink.table.types.logical.YearMonthIntervalType;
5455
import org.apache.flink.table.types.utils.DataTypeFactoryMock;
5556

@@ -152,7 +153,11 @@ private static Stream<LogicalType> testConversion() {
152153
new MultisetType(BinaryType.ofEmptyLiteral()),
153154
new MultisetType(VarBinaryType.ofEmptyLiteral()),
154155
new BitmapType(),
156+
new VariantType(),
157+
new ArrayType(new VariantType()),
158+
new MapType(new VarCharType(5), new VariantType()),
155159
RowType.of(new BigIntType(), new IntType(false), new VarCharType(200)),
160+
RowType.of(new VariantType(), new VariantType(false)),
156161
RowType.of(
157162
new LogicalType[] {
158163
new BigIntType(), new IntType(false), new VarCharType(200)

0 commit comments

Comments
 (0)