This is an automated email from the ASF dual-hosted git repository.

ferenc-csaky pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git


The following commit(s) were added to refs/heads/master by this push:
     new 0c576881a7d [FLINK-40376][formats] Support `VARIANT` type in 
`flink-json` converters
0c576881a7d is described below

commit 0c576881a7dea7619a6ce9bedb18c1b577e6c3ad
Author: Ferenc Csaky <[email protected]>
AuthorDate: Thu Aug 13 17:14:26 2026 +0200

    [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 +++
 .../formats/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 ea5ced4c905..097ed8c021d 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.MapType;
 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 @@ public class JsonParserToRowDataConverters implements 
Serializable {
             case BINARY:
             case VARBINARY:
                 return JsonParser::getBinaryValue;
+            case VARIANT:
+                return this::convertToVariant;
             case DECIMAL:
                 return createDecimalConverter((DecimalType) type);
             case ARRAY:
@@ -323,6 +327,10 @@ public class JsonParserToRowDataConverters implements 
Serializable {
         }
     }
 
+    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 653b8b3621b..dedc58e0bff 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.MapType;
 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 @@ public class JsonToRowDataConverters implements 
Serializable {
             case BINARY:
             case VARBINARY:
                 return this::convertToBytes;
+            case VARIANT:
+                return this::convertToVariant;
             case DECIMAL:
                 return createDecimalConverter((DecimalType) type);
             case ARRAY:
@@ -270,6 +274,14 @@ public class JsonToRowDataConverters implements 
Serializable {
         }
     }
 
+    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 11ab01a35e8..c43f20f92e3 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.LogicalTypeFamily;
 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 @@ public class RowDataToJsonConverters implements 
Serializable {
                         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 @@ public class RowDataToJsonConverters implements 
Serializable {
         };
     }
 
+    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 fc320bc1696..acc09afb7b9 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.ParameterizedTe
 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.TIME;
 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 @@ public class JsonRowDataSerDeSchemaTest {
         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<RowData> 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<RowData> 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<RowData> 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()));

Reply via email to