From 1360372965257224ba0cd3bfdf823ef849e97b4e Mon Sep 17 00:00:00 2001 From: Manu Zhang Date: Sun, 9 Aug 2026 23:19:56 +0800 Subject: [PATCH 1/4] Parquet: Validate geospatial projection parameters Generated-by: Codex --- .../iceberg/parquet/MessageTypeToType.java | 4 + .../apache/iceberg/parquet/PruneColumns.java | 12 +++ .../iceberg/parquet/TestPruneColumns.java | 93 +++++++++++++++++++ 3 files changed, 109 insertions(+) diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/MessageTypeToType.java b/parquet/src/main/java/org/apache/iceberg/parquet/MessageTypeToType.java index 2b01bf882e75..5d791e66f35c 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/MessageTypeToType.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/MessageTypeToType.java @@ -149,6 +149,10 @@ public Type map(GroupType map, Type keyType, Type valueType) { @Override public Type primitive(PrimitiveType primitive) { + return convertPrimitive(primitive); + } + + static Type convertPrimitive(PrimitiveType primitive) { // first, use the logical type annotation, if present LogicalTypeAnnotation logicalType = primitive.getLogicalTypeAnnotation(); if (logicalType != null) { diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java b/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java index 0647a09f53fe..5c94426d3956 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java @@ -18,12 +18,16 @@ */ package org.apache.iceberg.parquet; +import static org.apache.iceberg.types.Type.TypeID.GEOGRAPHY; +import static org.apache.iceberg.types.Type.TypeID.GEOMETRY; + import java.util.Collections; import java.util.List; import java.util.Set; import org.apache.iceberg.relocated.com.google.common.base.Objects; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.types.Types.ListType; import org.apache.iceberg.types.Types.MapType; import org.apache.iceberg.types.Types.NestedField; @@ -162,6 +166,14 @@ public Type variant( @Override public Type primitive( org.apache.iceberg.types.Type.PrimitiveType expected, PrimitiveType primitive) { + if (expected != null && (expected.typeId() == GEOMETRY || expected.typeId() == GEOGRAPHY)) { + Preconditions.checkArgument( + TypeUtil.isPromotionAllowed(MessageTypeToType.convertPrimitive(primitive), expected), + "Cannot read Parquet type %s as Iceberg type %s", + primitive, + expected); + } + return null; } diff --git a/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java b/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java index 619b2c5a3470..6ea482dff4a4 100644 --- a/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java +++ b/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java @@ -19,10 +19,14 @@ package org.apache.iceberg.parquet; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import org.apache.iceberg.Schema; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.types.EdgeAlgorithm; import org.apache.iceberg.types.Types.DoubleType; +import org.apache.iceberg.types.Types.GeographyType; +import org.apache.iceberg.types.Types.GeometryType; import org.apache.iceberg.types.Types.IntegerType; import org.apache.iceberg.types.Types.ListType; import org.apache.iceberg.types.Types.MapType; @@ -31,6 +35,7 @@ import org.apache.iceberg.types.Types.StructType; import org.apache.iceberg.types.Types.VariantType; import org.apache.iceberg.variants.Variant; +import org.apache.parquet.column.schema.EdgeInterpolationAlgorithm; import org.apache.parquet.schema.LogicalTypeAnnotation; import org.apache.parquet.schema.MessageType; import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName; @@ -305,6 +310,94 @@ public void testVariant() { assertThat(actual).as("Pruned schema should be matched").isEqualTo(expected); } + @Test + public void acceptsMatchingGeospatialParameters() { + MessageType fileSchema = + Types.buildMessage() + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.geometryType(null)) + .id(1) + .named("geom_default") + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.geometryType("EPSG:3857")) + .id(2) + .named("geom_projected") + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.geographyType(null, null)) + .id(3) + .named("geog_default") + .optional(PrimitiveTypeName.BINARY) + .as( + LogicalTypeAnnotation.geographyType( + "EPSG:4326", EdgeInterpolationAlgorithm.ANDOYER)) + .id(4) + .named("geog_custom") + .named("table"); + + Schema projection = + new Schema( + NestedField.optional(1, "geom_default", GeometryType.crs84()), + NestedField.optional(2, "geom_projected", GeometryType.of("epsg:3857")), + NestedField.optional(3, "geog_default", GeographyType.crs84()), + NestedField.optional( + 4, "geog_custom", GeographyType.of("epsg:4326", EdgeAlgorithm.ANDOYER))); + + assertThat(ParquetSchemaUtil.pruneColumns(fileSchema, projection)).isEqualTo(fileSchema); + } + + @Test + public void rejectsGeometryCrsMismatch() { + MessageType fileSchema = + Types.buildMessage() + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.geometryType("EPSG:3857")) + .id(1) + .named("geom") + .named("table"); + Schema projection = new Schema(NestedField.optional(1, "geom", GeometryType.crs84())); + + assertThatThrownBy(() -> ParquetSchemaUtil.pruneColumns(fileSchema, projection)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Cannot read Parquet type") + .hasMessageContaining("geometry(OGC:CRS84)"); + } + + @Test + public void rejectsGeographyCrsMismatch() { + MessageType fileSchema = + Types.buildMessage() + .optional(PrimitiveTypeName.BINARY) + .as( + LogicalTypeAnnotation.geographyType( + "EPSG:4326", EdgeInterpolationAlgorithm.SPHERICAL)) + .id(1) + .named("geog") + .named("table"); + Schema projection = new Schema(NestedField.optional(1, "geog", GeographyType.crs84())); + + assertThatThrownBy(() -> ParquetSchemaUtil.pruneColumns(fileSchema, projection)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Cannot read Parquet type") + .hasMessageContaining("geography(OGC:CRS84, spherical)"); + } + + @Test + public void rejectsGeographyAlgorithmMismatch() { + MessageType fileSchema = + Types.buildMessage() + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.geographyType("OGC:CRS84", EdgeInterpolationAlgorithm.KARNEY)) + .id(1) + .named("geog") + .named("table"); + Schema projection = new Schema(NestedField.optional(1, "geog", GeographyType.crs84())); + + assertThatThrownBy(() -> ParquetSchemaUtil.pruneColumns(fileSchema, projection)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Cannot read Parquet type") + .hasMessageContaining("geography(OGC:CRS84, spherical)"); + } + private static Type buildVariantType(int id, String name) { return Types.buildGroup(Type.Repetition.OPTIONAL) .as(LogicalTypeAnnotation.variantType(Variant.VARIANT_SPEC_VERSION)) From a6bc8b596c56e1c2039ef498b51030c28aef1086 Mon Sep 17 00:00:00 2001 From: Manu Zhang Date: Sun, 9 Aug 2026 23:29:02 +0800 Subject: [PATCH 2/4] Parquet: Validate fallback geospatial projections Generated-by: Codex --- .../iceberg/parquet/ParquetSchemaUtil.java | 6 ++++++ .../org/apache/iceberg/parquet/PruneColumns.java | 8 ++++++-- .../apache/iceberg/parquet/TestPruneColumns.java | 16 ++++++++++++++++ 3 files changed, 28 insertions(+), 2 deletions(-) diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetSchemaUtil.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetSchemaUtil.java index 9a81626827c6..a74fec8697c8 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetSchemaUtil.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetSchemaUtil.java @@ -159,6 +159,12 @@ public static MessageType pruneColumnsFallback(MessageType fileSchema, Schema ex int ordinal = 1; for (Type type : fileSchema.getFields()) { if (selectedIds.contains(ordinal)) { + Types.NestedField expectedField = expectedSchema.findField(ordinal); + if (type.isPrimitive() && expectedField.type().isPrimitiveType()) { + PruneColumns.validatePrimitive( + expectedField.type().asPrimitiveType(), type.asPrimitiveType()); + } + builder.addField(type.withId(ordinal)); } ordinal += 1; diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java b/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java index 5c94426d3956..affe67e69e6f 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java @@ -166,6 +166,12 @@ public Type variant( @Override public Type primitive( org.apache.iceberg.types.Type.PrimitiveType expected, PrimitiveType primitive) { + validatePrimitive(expected, primitive); + return null; + } + + static void validatePrimitive( + org.apache.iceberg.types.Type.PrimitiveType expected, PrimitiveType primitive) { if (expected != null && (expected.typeId() == GEOMETRY || expected.typeId() == GEOGRAPHY)) { Preconditions.checkArgument( TypeUtil.isPromotionAllowed(MessageTypeToType.convertPrimitive(primitive), expected), @@ -173,8 +179,6 @@ public Type primitive( primitive, expected); } - - return null; } private Integer getId(Type type) { diff --git a/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java b/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java index 6ea482dff4a4..c1120cfa50be 100644 --- a/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java +++ b/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java @@ -362,6 +362,22 @@ public void rejectsGeometryCrsMismatch() { .hasMessageContaining("geometry(OGC:CRS84)"); } + @Test + public void rejectsGeometryCrsMismatchWithoutIds() { + MessageType fileSchema = + Types.buildMessage() + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.geometryType("EPSG:3857")) + .named("geom") + .named("table"); + Schema projection = new Schema(NestedField.optional(1, "geom", GeometryType.crs84())); + + assertThatThrownBy(() -> ParquetSchemaUtil.pruneColumnsFallback(fileSchema, projection)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Cannot read Parquet type") + .hasMessageContaining("geometry(OGC:CRS84)"); + } + @Test public void rejectsGeographyCrsMismatch() { MessageType fileSchema = From d1e612fb31f6870b3606d863ee5b200a71c6614b Mon Sep 17 00:00:00 2001 From: Manu Zhang Date: Sun, 9 Aug 2026 23:32:12 +0800 Subject: [PATCH 3/4] Parquet: Fix geospatial validation Checkstyle Generated-by: Codex --- .../main/java/org/apache/iceberg/parquet/PruneColumns.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java b/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java index affe67e69e6f..3471374fce89 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/PruneColumns.java @@ -18,15 +18,13 @@ */ package org.apache.iceberg.parquet; -import static org.apache.iceberg.types.Type.TypeID.GEOGRAPHY; -import static org.apache.iceberg.types.Type.TypeID.GEOMETRY; - import java.util.Collections; import java.util.List; import java.util.Set; import org.apache.iceberg.relocated.com.google.common.base.Objects; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.types.Type.TypeID; import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.types.Types.ListType; import org.apache.iceberg.types.Types.MapType; @@ -172,7 +170,8 @@ public Type primitive( static void validatePrimitive( org.apache.iceberg.types.Type.PrimitiveType expected, PrimitiveType primitive) { - if (expected != null && (expected.typeId() == GEOMETRY || expected.typeId() == GEOGRAPHY)) { + if (expected != null + && (expected.typeId() == TypeID.GEOMETRY || expected.typeId() == TypeID.GEOGRAPHY)) { Preconditions.checkArgument( TypeUtil.isPromotionAllowed(MessageTypeToType.convertPrimitive(primitive), expected), "Cannot read Parquet type %s as Iceberg type %s", From e86bf5a4bfebfaa3db75a85a96952b26b172eb48 Mon Sep 17 00:00:00 2001 From: Manu Zhang Date: Mon, 10 Aug 2026 23:42:31 +0800 Subject: [PATCH 4/4] Parquet: Assert mismatched geography algorithm Generated-by: Codex --- .../test/java/org/apache/iceberg/parquet/TestPruneColumns.java | 1 + 1 file changed, 1 insertion(+) diff --git a/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java b/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java index c1120cfa50be..1224b2c1bfac 100644 --- a/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java +++ b/parquet/src/test/java/org/apache/iceberg/parquet/TestPruneColumns.java @@ -411,6 +411,7 @@ public void rejectsGeographyAlgorithmMismatch() { assertThatThrownBy(() -> ParquetSchemaUtil.pruneColumns(fileSchema, projection)) .isInstanceOf(IllegalArgumentException.class) .hasMessageContaining("Cannot read Parquet type") + .hasMessageContaining("KARNEY") .hasMessageContaining("geography(OGC:CRS84, spherical)"); }