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