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