From d3515c95322db6406c4b526c3ab4d575aba7a6a9 Mon Sep 17 00:00:00 2001 From: Ferenc Csaky Date: Wed, 12 Aug 2026 12:46:26 +0200 Subject: [PATCH] [FLINK-40376][formats] Support `VARIANT` type in `flink-json` converters Generated-by: openai/gpt-5.6-terra --- .../json/JsonParserToRowDataConverters.java | 8 ++ .../formats/json/JsonToRowDataConverters.java | 12 ++ .../formats/json/RowDataToJsonConverters.java | 12 ++ .../json/JsonRowDataSerDeSchemaTest.java | 109 ++++++++++++++++++ 4 files changed, 141 insertions(+) diff --git a/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonParserToRowDataConverters.java b/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonParserToRowDataConverters.java index ea5ced4c9058c3..097ed8c021daa8 100644 --- a/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonParserToRowDataConverters.java +++ b/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonParserToRowDataConverters.java @@ -37,6 +37,8 @@ import org.apache.flink.table.types.logical.MultisetType; import org.apache.flink.table.types.logical.RowType; import org.apache.flink.table.types.logical.utils.LogicalTypeUtils; +import org.apache.flink.types.variant.BinaryVariant; +import org.apache.flink.types.variant.BinaryVariantInternalBuilder; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonParser; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonToken; @@ -167,6 +169,8 @@ private JsonParserToRowDataConverter createNotNullConverter(LogicalType type) { case BINARY: case VARBINARY: return JsonParser::getBinaryValue; + case VARIANT: + return this::convertToVariant; case DECIMAL: return createDecimalConverter((DecimalType) type); case ARRAY: @@ -323,6 +327,10 @@ private StringData convertToString(JsonParser jp) throws IOException { } } + private BinaryVariant convertToVariant(JsonParser jp) throws IOException { + return BinaryVariantInternalBuilder.parseJson(jp.readValueAsTree().toString(), false); + } + private JsonParserToRowDataConverter createDecimalConverter(DecimalType decimalType) { final int precision = decimalType.getPrecision(); final int scale = decimalType.getScale(); diff --git a/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonToRowDataConverters.java b/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonToRowDataConverters.java index 653b8b3621b417..dedc58e0bff8e2 100644 --- a/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonToRowDataConverters.java +++ b/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/JsonToRowDataConverters.java @@ -37,6 +37,8 @@ import org.apache.flink.table.types.logical.MultisetType; import org.apache.flink.table.types.logical.RowType; import org.apache.flink.table.types.logical.utils.LogicalTypeUtils; +import org.apache.flink.types.variant.BinaryVariant; +import org.apache.flink.types.variant.BinaryVariantInternalBuilder; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ArrayNode; @@ -137,6 +139,8 @@ private JsonToRowDataConverter createNotNullConverter(LogicalType type) { case BINARY: case VARBINARY: return this::convertToBytes; + case VARIANT: + return this::convertToVariant; case DECIMAL: return createDecimalConverter((DecimalType) type); case ARRAY: @@ -270,6 +274,14 @@ private StringData convertToString(JsonNode jsonNode) { } } + private BinaryVariant convertToVariant(JsonNode jsonNode) { + try { + return BinaryVariantInternalBuilder.parseJson(jsonNode.toString(), false); + } catch (IOException e) { + throw new JsonParseException("Unable to deserialize VARIANT value.", e); + } + } + private byte[] convertToBytes(JsonNode jsonNode) { try { return jsonNode.binaryValue(); diff --git a/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/RowDataToJsonConverters.java b/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/RowDataToJsonConverters.java index 11ab01a35e8a40..c43f20f92e391e 100644 --- a/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/RowDataToJsonConverters.java +++ b/flink-formats/flink-json/src/main/java/org/apache/flink/formats/json/RowDataToJsonConverters.java @@ -33,12 +33,14 @@ import org.apache.flink.table.types.logical.MapType; import org.apache.flink.table.types.logical.MultisetType; import org.apache.flink.table.types.logical.RowType; +import org.apache.flink.types.variant.Variant; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ArrayNode; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ObjectNode; +import java.io.IOException; import java.io.Serializable; import java.math.BigDecimal; import java.time.LocalDate; @@ -149,6 +151,8 @@ private RowDataToJsonConverter createNotNullConverter(LogicalType type) { new IntType()); case ROW: return createRowConverter((RowType) type); + case VARIANT: + return this::convertVariant; case RAW: default: throw new UnsupportedOperationException("Not support to parse type: " + type); @@ -166,6 +170,14 @@ private RowDataToJsonConverter createDecimalConverter() { }; } + private JsonNode convertVariant(ObjectMapper mapper, JsonNode reuse, Object value) { + try { + return mapper.readTree(((Variant) value).toJson()); + } catch (IOException e) { + throw new JsonParseException("Unable to serialize VARIANT value.", e); + } + } + private RowDataToJsonConverter createDateConverter() { return (mapper, reuse, value) -> { int days = (int) value; diff --git a/flink-formats/flink-json/src/test/java/org/apache/flink/formats/json/JsonRowDataSerDeSchemaTest.java b/flink-formats/flink-json/src/test/java/org/apache/flink/formats/json/JsonRowDataSerDeSchemaTest.java index fc320bc16963c8..acc09afb7b9693 100644 --- a/flink-formats/flink-json/src/test/java/org/apache/flink/formats/json/JsonRowDataSerDeSchemaTest.java +++ b/flink-formats/flink-json/src/test/java/org/apache/flink/formats/json/JsonRowDataSerDeSchemaTest.java @@ -36,6 +36,7 @@ import org.apache.flink.testutils.junit.extensions.parameterized.Parameters; import org.apache.flink.testutils.logging.LoggerAuditingExtension; import org.apache.flink.types.Row; +import org.apache.flink.types.variant.BinaryVariantBuilder; import org.apache.flink.util.Collector; import org.apache.flink.util.jackson.JacksonMapperFactory; @@ -87,6 +88,7 @@ import static org.apache.flink.table.api.DataTypes.TIMESTAMP; import static org.apache.flink.table.api.DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE; import static org.apache.flink.table.api.DataTypes.TINYINT; +import static org.apache.flink.table.api.DataTypes.VARIANT; import static org.apache.flink.table.types.utils.TypeConversions.fromLogicalToDataType; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -239,6 +241,113 @@ void testSerDe() throws Exception { assertThat(serializedJson).containsExactly(actualBytes); } + @TestTemplate + void testVariantScalarsAndNestedTypes() throws Exception { + String json = + "{\"string\":\"Flink\",\"boolean\":true,\"integer\":42,\"decimal\":3.14," + + "\"array\":[1,{\"nested\":\"value\"}],\"nestedArray\":[\"array\"]," + + "\"nestedMap\":{\"key\":\"map\"},\"nestedRow\":{\"field\":\"row\"}}"; + DataType dataType = + ROW( + FIELD("string", VARIANT()), + FIELD("boolean", VARIANT()), + FIELD("integer", VARIANT()), + FIELD("decimal", VARIANT()), + FIELD("array", VARIANT()), + FIELD("nestedArray", ARRAY(VARIANT())), + FIELD("nestedMap", MAP(STRING(), VARIANT())), + FIELD("nestedRow", ROW(FIELD("field", VARIANT())))); + RowType rowType = (RowType) dataType.getLogicalType(); + + DeserializationSchema deserializationSchema = + createDeserializationSchema( + isJsonParser, rowType, false, false, TimestampFormat.ISO_8601); + open(deserializationSchema); + + RowData rowData = deserializationSchema.deserialize(json.getBytes()); + assertThat(rowData.getVariant(0).toJson()).isEqualTo("\"Flink\""); + assertThat(rowData.getVariant(1).toJson()).isEqualTo("true"); + assertThat(rowData.getVariant(2).toJson()).isEqualTo("42"); + assertThat(rowData.getVariant(3).toJson()).isEqualTo("3.14"); + assertThat(rowData.getVariant(4).toJson()).isEqualTo("[1,{\"nested\":\"value\"}]"); + assertThat(rowData.getArray(5).getVariant(0).toJson()).isEqualTo("\"array\""); + assertThat(rowData.getMap(6).valueArray().getVariant(0).toJson()).isEqualTo("\"map\""); + assertThat(rowData.getRow(7, 1).getVariant(0).toJson()).isEqualTo("\"row\""); + + JsonRowDataSerializationSchema serializationSchema = + new JsonRowDataSerializationSchema( + rowType, + TimestampFormat.ISO_8601, + JsonFormatOptions.MapNullKeyMode.LITERAL, + "null", + true, + false); + open(serializationSchema); + + assertThat(OBJECT_MAPPER.readTree(serializationSchema.serialize(rowData))) + .isEqualTo(OBJECT_MAPPER.readTree(json)); + } + + @TestTemplate + void testVariantDuplicateKeysUseLastValue() throws Exception { + byte[] json = "{\"variant\":{\"key\":1,\"key\":2}}".getBytes(); + RowType rowType = (RowType) ROW(FIELD("variant", VARIANT())).getLogicalType(); + + DeserializationSchema deserializationSchema = + createDeserializationSchema( + isJsonParser, rowType, false, false, TimestampFormat.ISO_8601); + open(deserializationSchema); + assertThat(deserializationSchema.deserialize(json).getVariant(0).toJson()) + .isEqualTo("{\"key\":2}"); + } + + @TestTemplate + void testVariantJsonNullIsDeserializedAsSqlNull() throws Exception { + byte[] json = "{\"variant\":null}".getBytes(); + RowType rowType = (RowType) ROW(FIELD("variant", VARIANT())).getLogicalType(); + + DeserializationSchema deserializationSchema = + createDeserializationSchema( + isJsonParser, rowType, false, false, TimestampFormat.ISO_8601); + open(deserializationSchema); + assertThat(deserializationSchema.deserialize(json).isNullAt(0)).isTrue(); + + JsonRowDataSerializationSchema serializationSchema = + new JsonRowDataSerializationSchema( + rowType, + TimestampFormat.ISO_8601, + JsonFormatOptions.MapNullKeyMode.LITERAL, + "null", + true, + false); + open(serializationSchema); + assertThat( + serializationSchema.serialize( + GenericRowData.of(new BinaryVariantBuilder().ofNull()))) + .containsExactly(json); + } + + @Test + void testVariantSerializationRejectsNonFiniteNumbers() { + RowType rowType = (RowType) ROW(FIELD("variant", VARIANT())).getLogicalType(); + JsonRowDataSerializationSchema serializationSchema = + new JsonRowDataSerializationSchema( + rowType, + TimestampFormat.ISO_8601, + JsonFormatOptions.MapNullKeyMode.LITERAL, + "null", + true, + false); + open(serializationSchema); + + assertThatThrownBy( + () -> + serializationSchema.serialize( + GenericRowData.of( + new BinaryVariantBuilder().of(Double.NaN)))) + .hasMessage("Non-finite value NaN cannot be serialized to JSON."); + } + @Test public void testEmptyJsonArrayDeserialization() throws Exception { DataType dataType = ROW(FIELD("f1", INT()), FIELD("f2", BOOLEAN()), FIELD("f3", STRING()));