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));
+ }
}
}