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

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


The following commit(s) were added to refs/heads/master by this push:
     new bbb71a0a475 Fix IcebergIO failing to write EnumerationType (proto 
enum) fields (#40300)
bbb71a0a475 is described below

commit bbb71a0a475cb2f431a4047c0ea1c6adf769944f
Author: Manan Mangal <[email protected]>
AuthorDate: Wed Sep 30 14:39:12 2026 -0700

    Fix IcebergIO failing to write EnumerationType (proto enum) fields (#40300)
    
    EnumerationType (used for proto enum fields) wasn't handled in
    IcebergUtils's Beam-to-Iceberg schema and row conversion. Writing a Row
    with such a field threw "Unsupported Beam logical type Enum" while
    deriving the table schema, and threw a ClassCastException converting
    the row's value even against a pre-existing table.
    
    Map EnumerationType to Iceberg's string type on write, the same way
    Avro and BigQuery already convert it.
    
    Fixes #40299
---
 CHANGES.md                                         |  1 +
 .../apache/beam/sdk/io/iceberg/IcebergUtils.java   | 11 ++++++-
 .../beam/sdk/io/iceberg/IcebergUtilsTest.java      | 38 ++++++++++++++++++++++
 3 files changed, 49 insertions(+), 1 deletion(-)

diff --git a/CHANGES.md b/CHANGES.md
index 336773b93e5..5ef4f34c92b 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -85,6 +85,7 @@
 * (Java) Fixed the declared schema of the error output of the Kafka write 
SchemaTransform, which wrapped the error schema a second time and did not match 
the rows it emits ([#39760](https://github.com/apache/beam/issues/39760)).
 * (Go) Fixed pubsubio importing a `google.golang.org/genproto` package removed 
in recent releases, which broke builds of Go modules depending on a current 
`genproto` version ([#40018](https://github.com/apache/beam/issues/40018)).
 * (Java) BigQueryIO now treats a 404 when deleting a temporary table or 
dataset as success, so a replayed work item whose earlier attempt already 
deleted it no longer retries forever 
([#24997](https://github.com/apache/beam/issues/24997)).
+* (Java) IcebergIO now writes rows containing `EnumerationType` (proto enum) 
fields as strings, instead of throwing `Unsupported Beam logical type Enum` 
([#40299](https://github.com/apache/beam/issues/40299)).
 * Fixed X (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
 
 ## Security Fixes
diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
index fa8d17f3c47..8db3b643cf8 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergUtils.java
@@ -36,6 +36,7 @@ import java.util.UUID;
 import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
 import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.schemas.logicaltypes.EnumerationType;
 import org.apache.beam.sdk.schemas.logicaltypes.FixedPrecisionNumeric;
 import org.apache.beam.sdk.schemas.logicaltypes.MicrosInstant;
 import org.apache.beam.sdk.schemas.logicaltypes.PassThroughLogicalType;
@@ -84,6 +85,7 @@ public class IcebergUtils {
           .put(SqlTypes.DATETIME.getIdentifier(), 
Types.TimestampType.withoutZone())
           .put(SqlTypes.UUID.getIdentifier(), Types.UUIDType.get())
           .put(MicrosInstant.IDENTIFIER, Types.TimestampType.withZone())
+          .put(EnumerationType.IDENTIFIER, Types.StringType.get())
           .build();
 
   private static Schema.FieldType icebergTypeToBeamFieldType(
@@ -389,7 +391,14 @@ public class IcebergUtils {
         rec.setField(name, getIcebergTimestampValue(val, 
ts.shouldAdjustToUTC()));
         break;
       case STRING:
-        Optional.ofNullable(value.getString(name)).ifPresent(v -> 
rec.setField(name, v));
+        Object beamValue = value.getValue(name);
+        if (beamValue instanceof EnumerationType.Value) { // EnumerationType
+          EnumerationType enumType =
+              
value.getSchema().getField(name).getType().getLogicalType(EnumerationType.class);
+          rec.setField(name, enumType.toString((EnumerationType.Value) 
beamValue));
+        } else { // FieldType.STRING, or a String-backed PassThroughLogicalType
+          Optional.ofNullable((String) beamValue).ifPresent(v -> 
rec.setField(name, v));
+        }
         break;
       case UUID:
         Optional.ofNullable(value.getBytes(name))
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
index 80e5e2195f9..aa8f6799a3e 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergUtilsTest.java
@@ -39,6 +39,7 @@ import java.util.Arrays;
 import java.util.List;
 import java.util.Map;
 import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.schemas.logicaltypes.EnumerationType;
 import org.apache.beam.sdk.schemas.logicaltypes.FixedPrecisionNumeric;
 import org.apache.beam.sdk.schemas.logicaltypes.FixedString;
 import org.apache.beam.sdk.schemas.logicaltypes.SqlTypes;
@@ -289,6 +290,17 @@ public class IcebergUtilsTest {
       checkRowValueToRecordValue(Schema.FieldType.DECIMAL, 
Types.DecimalType.of(6, 3), num);
     }
 
+    @Test
+    public void testEnumeration() {
+      EnumerationType enumType = EnumerationType.create("ONE", "TWO", "THREE");
+
+      checkRowValueToRecordValue(
+          Schema.FieldType.logicalType(enumType),
+          enumType.valueOf("TWO"),
+          Types.StringType.get(),
+          "TWO");
+    }
+
     @Test
     public void testStruct() {
       Schema schema = Schema.builder().addStringField("nested_str").build();
@@ -1151,5 +1163,31 @@ public class IcebergUtilsTest {
 
       
assertTrue(convertedIcebergSchema.sameSchema(ICEBERG_SCHEMA_JDBC_ALL_TYPES));
     }
+
+    static final EnumerationType TEST_ENUMERATION_TYPE =
+        EnumerationType.create("ONE", "TWO", "THREE");
+
+    static final Schema BEAM_SCHEMA_ENUMERATION =
+        Schema.builder()
+            .addLogicalTypeField("status", TEST_ENUMERATION_TYPE)
+            .addNullableField(
+                "optional_status", 
Schema.FieldType.logicalType(TEST_ENUMERATION_TYPE))
+            .build();
+
+    static final org.apache.iceberg.Schema ICEBERG_SCHEMA_ENUMERATION =
+        new org.apache.iceberg.Schema(
+            required(1, "status", Types.StringType.get()),
+            optional(2, "optional_status", Types.StringType.get()));
+
+    @Test
+    public void testEnumerationBeamSchemaToIcebergSchema() {
+      // Iceberg has no enum type, so EnumerationType fields are stored as 
strings. This is a
+      // one-way conversion: converting ICEBERG_SCHEMA_ENUMERATION back with
+      // icebergSchemaToBeamSchema will not recover EnumerationType, only 
plain STRING fields.
+      org.apache.iceberg.Schema convertedIcebergSchema =
+          IcebergUtils.beamSchemaToIcebergSchema(BEAM_SCHEMA_ENUMERATION);
+
+      
assertTrue(convertedIcebergSchema.sameSchema(ICEBERG_SCHEMA_ENUMERATION));
+    }
   }
 }

Reply via email to