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

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 39de52c52a [iceberg] Support the variant type in Iceberg schema type 
conversion (#8795)
39de52c52a is described below

commit 39de52c52ac7357beb6b2ba26b1f93e6c3c69fb5
Author: Jiajia Li <[email protected]>
AuthorDate: Sun Jul 26 15:58:52 2026 +0800

    [iceberg] Support the variant type in Iceberg schema type conversion (#8795)
---
 .../paimon/iceberg/IcebergCommitCallback.java      | 57 +++++++++++++++++++++-
 .../paimon/iceberg/metadata/IcebergDataField.java  |  7 ++-
 .../paimon/iceberg/IcebergCompatibilityTest.java   | 13 +++++
 .../iceberg/metadata/IcebergDataFieldTest.java     | 51 +++++++++++++++++++
 .../org/apache/paimon/rest/RESTCatalogTest.java    |  5 +-
 5 files changed, 130 insertions(+), 3 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java
 
b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java
index df82c7dccc..d8fa8457db 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/iceberg/IcebergCommitCallback.java
@@ -51,6 +51,7 @@ import org.apache.paimon.manifest.ManifestEntry;
 import org.apache.paimon.options.Options;
 import org.apache.paimon.partition.PartitionPredicate;
 import org.apache.paimon.schema.SchemaManager;
+import org.apache.paimon.schema.TableSchema;
 import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.table.sink.CommitCallback;
 import org.apache.paimon.table.sink.TagCallback;
@@ -60,7 +61,11 @@ import org.apache.paimon.table.source.RawFile;
 import org.apache.paimon.table.source.ScanMode;
 import org.apache.paimon.table.source.snapshot.SnapshotReader;
 import org.apache.paimon.tag.Tag;
+import org.apache.paimon.types.ArrayType;
+import org.apache.paimon.types.DataField;
 import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.MapType;
+import org.apache.paimon.types.MultisetType;
 import org.apache.paimon.types.RowType;
 import org.apache.paimon.utils.DataFilePathFactories;
 import org.apache.paimon.utils.FileStorePathFactory;
@@ -77,6 +82,7 @@ import javax.annotation.Nullable;
 import java.io.IOException;
 import java.io.UncheckedIOException;
 import java.util.ArrayList;
+import java.util.Collection;
 import java.util.Collections;
 import java.util.Comparator;
 import java.util.HashMap;
@@ -547,6 +553,49 @@ public class IcebergCommitCallback implements 
CommitCallback, TagCallback {
         return result;
     }
 
+    /** VARIANT needs Iceberg row lineage, which Paimon Iceberg compatibility 
cannot publish. */
+    static void checkVariantNotPublishable(RowType rowType) {
+        Collection<String> variantFields = new LinkedHashSet<>();
+        for (DataField field : rowType.getFields()) {
+            collectVariantFields(field.name(), field.type(), variantFields);
+        }
+        Preconditions.checkArgument(
+                variantFields.isEmpty(),
+                "Columns %s use the VARIANT type, which Paimon Iceberg 
compatibility cannot "
+                        + "publish: it is an Iceberg format-version-3 type 
that requires row "
+                        + "lineage.",
+                variantFields);
+    }
+
+    private static void collectVariantFields(
+            String path, DataType type, Collection<String> variantFields) {
+        switch (type.getTypeRoot()) {
+            case VARIANT:
+                variantFields.add(path + ": " + type.asSQLString());
+                break;
+            case ARRAY:
+                collectVariantFields(
+                        path + ".element", ((ArrayType) 
type).getElementType(), variantFields);
+                break;
+            case MULTISET:
+                collectVariantFields(
+                        path + ".element", ((MultisetType) 
type).getElementType(), variantFields);
+                break;
+            case MAP:
+                collectVariantFields(path + ".key", ((MapType) 
type).getKeyType(), variantFields);
+                collectVariantFields(
+                        path + ".value", ((MapType) type).getValueType(), 
variantFields);
+                break;
+            case ROW:
+                for (DataField field : ((RowType) type).getFields()) {
+                    collectVariantFields(path + "." + field.name(), 
field.type(), variantFields);
+                }
+                break;
+            default:
+                break;
+        }
+    }
+
     // 
-------------------------------------------------------------------------------------
     // Create metadata based on old ones
     // 
-------------------------------------------------------------------------------------
@@ -1512,7 +1561,13 @@ public class IcebergCommitCallback implements 
CommitCallback, TagCallback {
 
         private IcebergSchema get(long schemaId) {
             return schemas.computeIfAbsent(
-                    schemaId, id -> 
IcebergSchema.create(schemaManager.schema(id)));
+                    schemaId,
+                    id -> {
+                        TableSchema schema = schemaManager.schema(id);
+                        // backstop: reject variant on each schema as it is 
emitted
+                        checkVariantNotPublishable(schema.logicalRowType());
+                        return IcebergSchema.create(schema);
+                    });
         }
 
         private long getLatestSchemaId() {
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/iceberg/metadata/IcebergDataField.java
 
b/paimon-core/src/main/java/org/apache/paimon/iceberg/metadata/IcebergDataField.java
index 95bbed9982..9862ff7f90 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/iceberg/metadata/IcebergDataField.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/iceberg/metadata/IcebergDataField.java
@@ -38,6 +38,7 @@ import org.apache.paimon.types.TimeType;
 import org.apache.paimon.types.TimestampType;
 import org.apache.paimon.types.VarBinaryType;
 import org.apache.paimon.types.VarCharType;
+import org.apache.paimon.types.VariantType;
 import org.apache.paimon.utils.Preconditions;
 
 import 
org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonCreator;
@@ -193,11 +194,13 @@ public class IcebergDataField {
                         timestampLtzPrecision >= 3 && timestampLtzPrecision <= 
9,
                         "Paimon Iceberg compatibility only support timestamp 
type with precision from 3 to 9.");
                 return timestampLtzPrecision >= 7 ? "timestamptz_ns" : 
"timestamptz";
+            case VARIANT:
+                return "variant";
             case ARRAY:
                 ArrayType arrayType = (ArrayType) dataType;
                 return new IcebergListType(
                         SpecialFields.getArrayElementFieldId(fieldId, depth + 
1),
-                        !dataType.isNullable(),
+                        !arrayType.getElementType().isNullable(),
                         toTypeObject(arrayType.getElementType(), fieldId, 
depth + 1));
             case MAP:
                 MapType mapType = (MapType) dataType;
@@ -285,6 +288,8 @@ public class IcebergDataField {
                     return new TimestampType(!isRequired, 9);
                 case "timestamptz_ns": // iceberg v3 format
                     return new LocalZonedTimestampType(!isRequired, 9);
+                case "variant": // iceberg v3 format
+                    return new VariantType(!isRequired);
                 default:
                     throw new UnsupportedOperationException(
                             "Unsupported primitive data type: " + icebergType);
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java
index ceb0b90695..c8a1b3fa4c 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/iceberg/IcebergCompatibilityTest.java
@@ -1790,4 +1790,17 @@ public class IcebergCompatibilityTest {
             }
         }
     }
+
+    @Test
+    public void testVariantIsNotPublishableToIceberg() {
+        RowType withVariant =
+                RowType.of(
+                        new DataType[] {DataTypes.INT(), DataTypes.VARIANT()},
+                        new String[] {"k", "payload"});
+        assertThatThrownBy(() -> 
IcebergCommitCallback.checkVariantNotPublishable(withVariant))
+                .hasMessageContaining("VARIANT type");
+
+        IcebergCommitCallback.checkVariantNotPublishable(
+                RowType.of(new DataType[] {DataTypes.INT()}, new String[] 
{"k"}));
+    }
 }
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/iceberg/metadata/IcebergDataFieldTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/iceberg/metadata/IcebergDataFieldTest.java
index 926a83b24e..0cfd5f0fe5 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/iceberg/metadata/IcebergDataFieldTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/iceberg/metadata/IcebergDataFieldTest.java
@@ -37,6 +37,7 @@ import org.apache.paimon.types.RowType;
 import org.apache.paimon.types.TimestampType;
 import org.apache.paimon.types.VarBinaryType;
 import org.apache.paimon.types.VarCharType;
+import org.apache.paimon.types.VariantType;
 import org.apache.paimon.utils.JsonSerdeUtil;
 
 import org.junit.jupiter.api.DisplayName;
@@ -269,6 +270,19 @@ class IcebergDataFieldTest {
         assertThat(listType.elementRequired()).isTrue();
     }
 
+    @Test
+    @DisplayName("Test array element-required follows the element, not the 
array")
+    void testArrayElementNullabilityIndependentOfArray() {
+        DataField arrayField =
+                new DataField(1, "array", new ArrayType(true, new 
VariantType(false)));
+        IcebergDataField icebergArray = new IcebergDataField(arrayField);
+
+        IcebergListType listType = (IcebergListType) icebergArray.type();
+        assertThat(listType.element()).isEqualTo("variant");
+        assertThat(listType.elementRequired()).isTrue();
+        assertThat(icebergArray.dataType()).isEqualTo(new ArrayType(true, new 
VariantType(false)));
+    }
+
     @Test
     @DisplayName("Test map type conversion")
     void testMapTypeConversion() {
@@ -461,6 +475,43 @@ class IcebergDataFieldTest {
         assertThat(timestamptzNsField.dataType()).isEqualTo(new 
LocalZonedTimestampType(true, 9));
     }
 
+    @Test
+    @DisplayName("Test variant type conversion")
+    void testVariantTypeConversion() {
+        DataField variantField = new DataField(1, "variant", new 
VariantType(true));
+        IcebergDataField icebergVariant = new IcebergDataField(variantField);
+
+        assertThat(icebergVariant.type()).isEqualTo("variant");
+        assertThat(icebergVariant.required()).isFalse();
+        assertThat(icebergVariant.dataType()).isEqualTo(new VariantType(true));
+    }
+
+    @Test
+    @DisplayName("Test variant type parsing")
+    void testVariantTypeParsing() {
+        IcebergDataField optionalVariant =
+                new IcebergDataField(1, "variant", false, "variant", null, 
"doc");
+        assertThat(optionalVariant.dataType()).isEqualTo(new 
VariantType(true));
+
+        IcebergDataField requiredVariant =
+                new IcebergDataField(2, "variant", true, "variant", null, 
"doc");
+        assertThat(requiredVariant.dataType()).isEqualTo(new 
VariantType(false));
+    }
+
+    @Test
+    @DisplayName("Test variant type serialization round trip")
+    void testVariantTypeSerializationRoundTrip() {
+        DataField variantField = new DataField(1, "variant", new 
VariantType(true));
+        IcebergDataField originalField = new IcebergDataField(variantField);
+
+        String json = JsonSerdeUtil.toJson(originalField);
+        IcebergDataField deserializedField = JsonSerdeUtil.fromJson(json, 
IcebergDataField.class);
+
+        assertThat(deserializedField.type()).isEqualTo("variant");
+        assertThat(deserializedField.dataType()).isEqualTo(new 
VariantType(true));
+        assertThat(deserializedField.toDatafield()).isEqualTo(variantField);
+    }
+
     @Test
     @DisplayName("Test unsupported primitive type parsing")
     void testUnsupportedPrimitiveTypeParsing() {
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/rest/RESTCatalogTest.java 
b/paimon-core/src/test/java/org/apache/paimon/rest/RESTCatalogTest.java
index 05376d3585..1438b49f17 100644
--- a/paimon-core/src/test/java/org/apache/paimon/rest/RESTCatalogTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/rest/RESTCatalogTest.java
@@ -3294,7 +3294,8 @@ public abstract class RESTCatalogTest extends 
CatalogTestBase {
                         Lists.newArrayList(
                                 new DataField(0, "pt", DataTypes.INT()),
                                 new DataField(1, "col1", DataTypes.STRING()),
-                                new DataField(2, "col2", DataTypes.STRING())),
+                                new DataField(2, "col2", DataTypes.STRING()),
+                                new DataField(3, "payload", 
DataTypes.VARIANT())),
                         Collections.singletonList("pt"),
                         Collections.emptyList(),
                         options,
@@ -3310,6 +3311,8 @@ public abstract class RESTCatalogTest extends 
CatalogTestBase {
         assertThat(tables).containsExactlyInAnyOrder("table1");
         assertThat(table.uuid()).isNotEmpty();
         assertThat(table.uuid()).isNotEqualTo(table.fullName());
+        assertThat(table.rowType().getField("payload").type())
+                .isInstanceOf(org.apache.paimon.types.VariantType.class);
     }
 
     @Test

Reply via email to