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()));