From f9d1338eae1c3b644b4c313017e6da80883a3995 Mon Sep 17 00:00:00 2001 From: rahulsmahadev Date: Fri, 4 Sep 2026 20:59:17 +0000 Subject: [PATCH] Spark 4.2: Set Iceberg field IDs as Spark column IDs for DSv2 validation Set StructField.withId with each Iceberg field ID during Schema to StructType conversion, so SparkTable.columns() exposes stable Column.id() for top-level and nested struct fields and Spark 4.2's nested-field column-ID validation (SPARK-57544) works with Iceberg tables. Adds a TestSparkSchemaUtil test. --- .../apache/iceberg/spark/TypeToSparkType.java | 4 ++- .../iceberg/spark/TestSparkSchemaUtil.java | 30 +++++++++++++++++++ 2 files changed, 33 insertions(+), 1 deletion(-) diff --git a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/TypeToSparkType.java b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/TypeToSparkType.java index dc077937577c..7290099d5b75 100644 --- a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/TypeToSparkType.java +++ b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/TypeToSparkType.java @@ -74,7 +74,9 @@ public DataType struct(Types.StructType struct, List fieldResults) { Types.NestedField field = fields.get(i); DataType type = fieldResults.get(i); Metadata metadata = fieldMetadata(field.fieldId()); - StructField sparkField = StructField.apply(field.name(), type, field.isOptional(), metadata); + StructField sparkField = + StructField.apply(field.name(), type, field.isOptional(), metadata) + .withId(Integer.toString(field.fieldId())); if (field.doc() != null) { sparkField = sparkField.withComment(field.doc()); } diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSchemaUtil.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSchemaUtil.java index ba0dfd84110b..d1c01d4e4d10 100644 --- a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSchemaUtil.java +++ b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSchemaUtil.java @@ -37,6 +37,8 @@ import org.apache.spark.sql.catalyst.expressions.MetadataAttribute; import org.apache.spark.sql.catalyst.types.DataTypeUtils; import org.apache.spark.sql.catalyst.util.ResolveDefaultColumnsUtils$; +import org.apache.spark.sql.connector.catalog.CatalogV2Util; +import org.apache.spark.sql.connector.catalog.Column; import org.apache.spark.sql.types.DataType; import org.apache.spark.sql.types.DataTypes; import org.apache.spark.sql.types.GeographyType; @@ -101,6 +103,34 @@ public void testSchemaConversionWithMetaDataColumnSchema() { } } + @Test + public void testFieldIDsInSparkColumns() { + Schema schema = + new Schema( + optional( + 10, + "location", + Types.StructType.of( + optional(11, "latitude", Types.DoubleType.get()), + optional(12, "longitude", Types.DoubleType.get())))); + + StructType sparkSchema = SparkSchemaUtil.convert(schema); + StructField location = sparkSchema.apply("location"); + StructType locationType = (StructType) location.dataType(); + + assertThat(location.id().get()).isEqualTo("10"); + assertThat(locationType.apply("latitude").id().get()).isEqualTo("11"); + assertThat(locationType.apply("longitude").id().get()).isEqualTo("12"); + + Column[] columns = CatalogV2Util.structTypeToV2Columns(sparkSchema, true); + assertThat(columns).hasSize(1); + assertThat(columns[0].id()).isEqualTo("10"); + + StructType columnType = (StructType) columns[0].dataType(); + assertThat(columnType.apply("latitude").id().get()).isEqualTo("11"); + assertThat(columnType.apply("longitude").id().get()).isEqualTo("12"); + } + @Test public void testGeospatialTypeConversion() { // a default-CRS geometry round-trips through the null <-> OGC:CRS84 normalization