NoahKusaba commented on code in PR #19:
URL: https://github.com/apache/datafusion-iceberg/pull/19#discussion_r4128890243


##########
crates/datafusion/src/physical_plan/project.rs:
##########
@@ -675,215 +688,142 @@ mod tests {
         assert_eq!(city_partition.value(1), "Los Angeles");
     }
 
-    #[test]
-    fn test_schema_validation_matching_schemas() {
-        use iceberg::TableIdent;
-        use iceberg::io::FileIO;
-        use iceberg::spec::{FormatVersion, NestedField, PrimitiveType, Schema, 
Type};
-
-        let table_schema = Arc::new(
-            Schema::builder()
-                .with_fields(vec![
-                    NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int))
-                        .into(),
-                    NestedField::required(
-                        2,
-                        "name",
-                        Type::Primitive(PrimitiveType::String),
-                    )
-                    .into(),
-                ])
-                .build()
-                .unwrap(),
-        );
-
+    /// A table with `fields`, partitioned by identity on its `id` column.
+    fn table_partitioned_by_id(fields: Vec<NestedField>) -> Table {
+        let table_schema = Schema::builder()
+            .with_fields(fields.into_iter().map(Arc::new))
+            .build()
+            .unwrap();
         let partition_spec = PartitionSpec::builder(table_schema.clone())
             .add_partition_field("id", "id_partition", Transform::Identity)
             .unwrap()
             .build()
             .unwrap();
-
-        let sort_order = iceberg::spec::SortOrder::builder()
-            .build(&table_schema)
-            .unwrap();
-
-        let table_metadata_builder = iceberg::spec::TableMetadataBuilder::new(
-            (*table_schema).clone(),
+        let sort_order = SortOrder::builder().build(&table_schema).unwrap();
+        let metadata = TableMetadataBuilder::new(
+            table_schema,
             partition_spec,
             sort_order,
             "/test/table".to_string(),
             FormatVersion::V2,
-            std::collections::HashMap::new(),
+            HashMap::new(),
         )
-        .unwrap();
-
-        let table_metadata = table_metadata_builder.build().unwrap();
-
-        // Create Arrow schema matching the table schema
-        let arrow_schema = Arc::new(ArrowSchema::new(vec![
-            Field::new("id", DataType::Int32, false),
-            Field::new("name", DataType::Utf8, false),
-        ]));
-
-        let input = Arc::new(EmptyExec::new(arrow_schema));
-
-        let table = Table::builder()
-            .metadata(table_metadata.metadata)
+        .unwrap()
+        .build()
+        .unwrap()
+        .metadata;
+        Table::builder()
+            .metadata(metadata)
             .identifier(TableIdent::from_strs(["test", "table"]).unwrap())
             .file_io(FileIO::new_with_fs())
             .metadata_location("/test/metadata.json")
             .runtime(test_runtime())
             .build()
-            .unwrap();
-
-        let result = project_with_partition(input, &table);
-        assert!(result.is_ok(), "Schema validation should pass");
-    }
-
-    #[test]
-    fn test_schema_validation_mismatched_schemas() {
-        use iceberg::TableIdent;
-        use iceberg::io::FileIO;
-        use iceberg::spec::{FormatVersion, NestedField, PrimitiveType, Schema, 
Type};
-
-        let table_schema = Arc::new(
-            Schema::builder()
-                .with_fields(vec![
-                    NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int))
-                        .into(),
-                    NestedField::required(
-                        2,
-                        "name",
-                        Type::Primitive(PrimitiveType::String),
-                    )
-                    .into(),
-                ])
-                .build()
-                .unwrap(),
-        );
-
-        let partition_spec = PartitionSpec::builder(table_schema.clone())
-            .add_partition_field("id", "id_partition", Transform::Identity)
             .unwrap()
-            .build()
-            .unwrap();
+    }
 
-        let sort_order = iceberg::spec::SortOrder::builder()
-            .build(&table_schema)
-            .unwrap();
+    /// Required `id: int` and `name: string` columns.
+    fn id_and_name_table() -> Table {
+        table_partitioned_by_id(vec![
+            NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)),
+            NestedField::required(2, "name", 
Type::Primitive(PrimitiveType::String)),
+        ])
+    }
 
-        let table_metadata_builder = iceberg::spec::TableMetadataBuilder::new(
-            (*table_schema).clone(),
-            partition_spec,
-            sort_order,
-            "/test/table".to_string(),
-            FormatVersion::V2,
-            std::collections::HashMap::new(),
-        )
-        .unwrap();
+    fn input_of(fields: Vec<Field>) -> Arc<dyn ExecutionPlan> {
+        Arc::new(EmptyExec::new(Arc::new(ArrowSchema::new(fields))))
+    }
 
-        let table_metadata = table_metadata_builder.build().unwrap();
+    const INCOMPATIBLE: &str = "Input schema is not compatible with Iceberg 
table schema";
 
-        // Create Arrow schema with different field name (mismatched)
-        let arrow_schema = Arc::new(ArrowSchema::new(vec![
+    #[test]
+    fn test_schema_validation_matching_schemas() {
+        let input = input_of(vec![
             Field::new("id", DataType::Int32, false),
-            Field::new("different_name", DataType::Utf8, false), // Wrong 
field name
-        ]));
+            Field::new("name", DataType::Utf8, false),
+        ]);
+        assert!(project_with_partition(input, &id_and_name_table()).is_ok());
+    }
 
-        let input = Arc::new(EmptyExec::new(arrow_schema));
+    #[test]
+    fn test_schema_validation_mismatched_schemas() {
+        let input = input_of(vec![
+            Field::new("id", DataType::Int32, false),
+            Field::new("different_name", DataType::Utf8, false),
+        ]);
+        let err = project_with_partition(input, &id_and_name_table())
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
+    }
 
-        let table = Table::builder()
-            .metadata(table_metadata.metadata)
-            .identifier(TableIdent::from_strs(["test", "table"]).unwrap())
-            .file_io(FileIO::new_with_fs())
-            .metadata_location("/test/metadata.json")
-            .runtime(test_runtime())
-            .build()
-            .unwrap();
+    #[test]
+    fn test_schema_validation_nullability() {
+        let id = |nullable| input_of(vec![Field::new("id", DataType::Int32, 
nullable)]);
+        let int = Type::Primitive(PrimitiveType::Int);
+
+        // A non-nullable input fits an optional column, e.g. an INSERT from a
+        // NOT NULL source.
+        let optional =
+            table_partitioned_by_id(vec![NestedField::optional(1, "id", 
int.clone())]);
+        assert!(project_with_partition(id(false), &optional).is_ok());
+        assert!(project_with_partition(id(true), &optional).is_ok());
+
+        // A nullable input could write nulls into a required column.
+        let required = table_partitioned_by_id(vec![NestedField::required(1, 
"id", int)]);
+        assert!(project_with_partition(id(false), &required).is_ok());
+        let err = project_with_partition(id(true), &required)
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
+    }
 
-        let result = project_with_partition(input, &table);
-        assert!(
-            result.is_err(),
-            "Schema validation should fail for mismatched schemas"
-        );
-        assert!(
-            result
-                .unwrap_err()
-                .to_string()
-                .contains("Input schema does not match Iceberg table schema")
-        );
+    #[test]
+    fn test_schema_validation_nested_nullability() {
+        let child = |nullable| Field::new("x", DataType::Int32, nullable);
+        let input = |nullable| {
+            input_of(vec![
+                Field::new("id", DataType::Int32, false),
+                Field::new(
+                    "s",
+                    DataType::Struct(Fields::from(vec![child(nullable)])),
+                    false,
+                ),
+            ])
+        };
+        let table = |x: NestedField| {
+            table_partitioned_by_id(vec![
+                NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)),
+                NestedField::required(
+                    2,
+                    "s",
+                    Type::Struct(StructType::new(vec![Arc::new(x)])),
+                ),
+            ])
+        };
+        let int = Type::Primitive(PrimitiveType::Int);
+
+        // The same rule applies inside a struct.
+        let optional = table(NestedField::optional(3, "x", int.clone()));
+        assert!(project_with_partition(input(false), &optional).is_ok());
+        assert!(project_with_partition(input(true), &optional).is_ok());
+
+        let required = table(NestedField::required(3, "x", int));
+        assert!(project_with_partition(input(false), &required).is_ok());
+        let err = project_with_partition(input(true), &required)
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
     }

Review Comment:
   Thanks, good call. Adding the list and map cases turned up a bug already on 
`main`: iceberg-rust's `strip_metadata_from_schema` fails on any list or map 
column with "Field stack underflow in list", so a partitioned INSERT with such 
a column errors out before the schemas are compared. `MetadataStripVisitor` 
only pushes onto its field stack in `before_field`, which the visitor doesn't 
call for list elements or map keys and values. Filed as 
apache/iceberg-rust#3297.
   
   The fix and a regression test are ready. I'll open that PR once my pending 
iceberg-rust PRs are merged (apache/iceberg-rust#3286, 
apache/iceberg-rust#2904), since reviews there are bandwidth-limited.
   
   In abee042 the four cases live in a shared `assert_nested_nullability` 
helper used by the struct test, and the list and map tests are [staged but 
commented 
out](https://github.com/apache/datafusion-iceberg/blob/abee0425e1c7af8ff0f5b323209c6d36e27fa4d1/crates/datafusion/src/physical_plan/project.rs#L831-L878)
 until the fix lands. They pass against `665c64e` with the fix applied, 
including the map case with `keys_sorted = false` (what Iceberg maps convert 
to), so enabling them is just uncommenting.
   
   Two list/map limitations of `contains` I noticed while writing them, neither 
new in this PR (the old `==` had both):
   - Field names are compared, so a list whose element is named `item` (Arrow's 
default, used by `DataType::new_list`) won't match Iceberg's `element`.
   - A source map with `keys_sorted = true` is rejected, since Iceberg maps 
convert to unsorted Arrow maps.
   
   If you think either is worth handling, I can fold them into #22 or open a 
separate issue.
   



##########
crates/datafusion/src/physical_plan/project.rs:
##########
@@ -675,215 +688,142 @@ mod tests {
         assert_eq!(city_partition.value(1), "Los Angeles");
     }
 
-    #[test]
-    fn test_schema_validation_matching_schemas() {
-        use iceberg::TableIdent;
-        use iceberg::io::FileIO;
-        use iceberg::spec::{FormatVersion, NestedField, PrimitiveType, Schema, 
Type};
-
-        let table_schema = Arc::new(
-            Schema::builder()
-                .with_fields(vec![
-                    NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int))
-                        .into(),
-                    NestedField::required(
-                        2,
-                        "name",
-                        Type::Primitive(PrimitiveType::String),
-                    )
-                    .into(),
-                ])
-                .build()
-                .unwrap(),
-        );
-
+    /// A table with `fields`, partitioned by identity on its `id` column.
+    fn table_partitioned_by_id(fields: Vec<NestedField>) -> Table {
+        let table_schema = Schema::builder()
+            .with_fields(fields.into_iter().map(Arc::new))
+            .build()
+            .unwrap();
         let partition_spec = PartitionSpec::builder(table_schema.clone())
             .add_partition_field("id", "id_partition", Transform::Identity)
             .unwrap()
             .build()
             .unwrap();
-
-        let sort_order = iceberg::spec::SortOrder::builder()
-            .build(&table_schema)
-            .unwrap();
-
-        let table_metadata_builder = iceberg::spec::TableMetadataBuilder::new(
-            (*table_schema).clone(),
+        let sort_order = SortOrder::builder().build(&table_schema).unwrap();
+        let metadata = TableMetadataBuilder::new(
+            table_schema,
             partition_spec,
             sort_order,
             "/test/table".to_string(),
             FormatVersion::V2,
-            std::collections::HashMap::new(),
+            HashMap::new(),
         )
-        .unwrap();
-
-        let table_metadata = table_metadata_builder.build().unwrap();
-
-        // Create Arrow schema matching the table schema
-        let arrow_schema = Arc::new(ArrowSchema::new(vec![
-            Field::new("id", DataType::Int32, false),
-            Field::new("name", DataType::Utf8, false),
-        ]));
-
-        let input = Arc::new(EmptyExec::new(arrow_schema));
-
-        let table = Table::builder()
-            .metadata(table_metadata.metadata)
+        .unwrap()
+        .build()
+        .unwrap()
+        .metadata;
+        Table::builder()
+            .metadata(metadata)
             .identifier(TableIdent::from_strs(["test", "table"]).unwrap())
             .file_io(FileIO::new_with_fs())
             .metadata_location("/test/metadata.json")
             .runtime(test_runtime())
             .build()
-            .unwrap();
-
-        let result = project_with_partition(input, &table);
-        assert!(result.is_ok(), "Schema validation should pass");
-    }
-
-    #[test]
-    fn test_schema_validation_mismatched_schemas() {
-        use iceberg::TableIdent;
-        use iceberg::io::FileIO;
-        use iceberg::spec::{FormatVersion, NestedField, PrimitiveType, Schema, 
Type};
-
-        let table_schema = Arc::new(
-            Schema::builder()
-                .with_fields(vec![
-                    NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int))
-                        .into(),
-                    NestedField::required(
-                        2,
-                        "name",
-                        Type::Primitive(PrimitiveType::String),
-                    )
-                    .into(),
-                ])
-                .build()
-                .unwrap(),
-        );
-
-        let partition_spec = PartitionSpec::builder(table_schema.clone())
-            .add_partition_field("id", "id_partition", Transform::Identity)
             .unwrap()
-            .build()
-            .unwrap();
+    }
 
-        let sort_order = iceberg::spec::SortOrder::builder()
-            .build(&table_schema)
-            .unwrap();
+    /// Required `id: int` and `name: string` columns.
+    fn id_and_name_table() -> Table {
+        table_partitioned_by_id(vec![
+            NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)),
+            NestedField::required(2, "name", 
Type::Primitive(PrimitiveType::String)),
+        ])
+    }
 
-        let table_metadata_builder = iceberg::spec::TableMetadataBuilder::new(
-            (*table_schema).clone(),
-            partition_spec,
-            sort_order,
-            "/test/table".to_string(),
-            FormatVersion::V2,
-            std::collections::HashMap::new(),
-        )
-        .unwrap();
+    fn input_of(fields: Vec<Field>) -> Arc<dyn ExecutionPlan> {
+        Arc::new(EmptyExec::new(Arc::new(ArrowSchema::new(fields))))
+    }
 
-        let table_metadata = table_metadata_builder.build().unwrap();
+    const INCOMPATIBLE: &str = "Input schema is not compatible with Iceberg 
table schema";
 
-        // Create Arrow schema with different field name (mismatched)
-        let arrow_schema = Arc::new(ArrowSchema::new(vec![
+    #[test]
+    fn test_schema_validation_matching_schemas() {
+        let input = input_of(vec![
             Field::new("id", DataType::Int32, false),
-            Field::new("different_name", DataType::Utf8, false), // Wrong 
field name
-        ]));
+            Field::new("name", DataType::Utf8, false),
+        ]);
+        assert!(project_with_partition(input, &id_and_name_table()).is_ok());
+    }
 
-        let input = Arc::new(EmptyExec::new(arrow_schema));
+    #[test]
+    fn test_schema_validation_mismatched_schemas() {
+        let input = input_of(vec![
+            Field::new("id", DataType::Int32, false),
+            Field::new("different_name", DataType::Utf8, false),
+        ]);
+        let err = project_with_partition(input, &id_and_name_table())
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
+    }
 
-        let table = Table::builder()
-            .metadata(table_metadata.metadata)
-            .identifier(TableIdent::from_strs(["test", "table"]).unwrap())
-            .file_io(FileIO::new_with_fs())
-            .metadata_location("/test/metadata.json")
-            .runtime(test_runtime())
-            .build()
-            .unwrap();
+    #[test]
+    fn test_schema_validation_nullability() {
+        let id = |nullable| input_of(vec![Field::new("id", DataType::Int32, 
nullable)]);
+        let int = Type::Primitive(PrimitiveType::Int);
+
+        // A non-nullable input fits an optional column, e.g. an INSERT from a
+        // NOT NULL source.
+        let optional =
+            table_partitioned_by_id(vec![NestedField::optional(1, "id", 
int.clone())]);
+        assert!(project_with_partition(id(false), &optional).is_ok());
+        assert!(project_with_partition(id(true), &optional).is_ok());
+
+        // A nullable input could write nulls into a required column.
+        let required = table_partitioned_by_id(vec![NestedField::required(1, 
"id", int)]);
+        assert!(project_with_partition(id(false), &required).is_ok());
+        let err = project_with_partition(id(true), &required)
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
+    }
 
-        let result = project_with_partition(input, &table);
-        assert!(
-            result.is_err(),
-            "Schema validation should fail for mismatched schemas"
-        );
-        assert!(
-            result
-                .unwrap_err()
-                .to_string()
-                .contains("Input schema does not match Iceberg table schema")
-        );
+    #[test]
+    fn test_schema_validation_nested_nullability() {
+        let child = |nullable| Field::new("x", DataType::Int32, nullable);
+        let input = |nullable| {
+            input_of(vec![
+                Field::new("id", DataType::Int32, false),
+                Field::new(
+                    "s",
+                    DataType::Struct(Fields::from(vec![child(nullable)])),
+                    false,
+                ),
+            ])
+        };
+        let table = |x: NestedField| {
+            table_partitioned_by_id(vec![
+                NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)),
+                NestedField::required(
+                    2,
+                    "s",
+                    Type::Struct(StructType::new(vec![Arc::new(x)])),
+                ),
+            ])
+        };
+        let int = Type::Primitive(PrimitiveType::Int);
+
+        // The same rule applies inside a struct.
+        let optional = table(NestedField::optional(3, "x", int.clone()));
+        assert!(project_with_partition(input(false), &optional).is_ok());
+        assert!(project_with_partition(input(true), &optional).is_ok());
+
+        let required = table(NestedField::required(3, "x", int));
+        assert!(project_with_partition(input(false), &required).is_ok());
+        let err = project_with_partition(input(true), &required)
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
     }

Review Comment:
   I'm going to wait until my other pending datafusion related PR's are merged 
into Iceberg-rust, as they are severely bandwidth restricted on reviews:
   - https://github.com/apache/iceberg-rust/pull/3286
   - https://github.com/apache/iceberg-rust/pull/2904
   
   I also have an additional 5 unrelated PR's I want to eventually merge into 
iceberg-rust that I have written down as well. 
   



##########
crates/datafusion/src/physical_plan/project.rs:
##########
@@ -675,215 +688,142 @@ mod tests {
         assert_eq!(city_partition.value(1), "Los Angeles");
     }
 
-    #[test]
-    fn test_schema_validation_matching_schemas() {
-        use iceberg::TableIdent;
-        use iceberg::io::FileIO;
-        use iceberg::spec::{FormatVersion, NestedField, PrimitiveType, Schema, 
Type};
-
-        let table_schema = Arc::new(
-            Schema::builder()
-                .with_fields(vec![
-                    NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int))
-                        .into(),
-                    NestedField::required(
-                        2,
-                        "name",
-                        Type::Primitive(PrimitiveType::String),
-                    )
-                    .into(),
-                ])
-                .build()
-                .unwrap(),
-        );
-
+    /// A table with `fields`, partitioned by identity on its `id` column.
+    fn table_partitioned_by_id(fields: Vec<NestedField>) -> Table {
+        let table_schema = Schema::builder()
+            .with_fields(fields.into_iter().map(Arc::new))
+            .build()
+            .unwrap();
         let partition_spec = PartitionSpec::builder(table_schema.clone())
             .add_partition_field("id", "id_partition", Transform::Identity)
             .unwrap()
             .build()
             .unwrap();
-
-        let sort_order = iceberg::spec::SortOrder::builder()
-            .build(&table_schema)
-            .unwrap();
-
-        let table_metadata_builder = iceberg::spec::TableMetadataBuilder::new(
-            (*table_schema).clone(),
+        let sort_order = SortOrder::builder().build(&table_schema).unwrap();
+        let metadata = TableMetadataBuilder::new(
+            table_schema,
             partition_spec,
             sort_order,
             "/test/table".to_string(),
             FormatVersion::V2,
-            std::collections::HashMap::new(),
+            HashMap::new(),
         )
-        .unwrap();
-
-        let table_metadata = table_metadata_builder.build().unwrap();
-
-        // Create Arrow schema matching the table schema
-        let arrow_schema = Arc::new(ArrowSchema::new(vec![
-            Field::new("id", DataType::Int32, false),
-            Field::new("name", DataType::Utf8, false),
-        ]));
-
-        let input = Arc::new(EmptyExec::new(arrow_schema));
-
-        let table = Table::builder()
-            .metadata(table_metadata.metadata)
+        .unwrap()
+        .build()
+        .unwrap()
+        .metadata;
+        Table::builder()
+            .metadata(metadata)
             .identifier(TableIdent::from_strs(["test", "table"]).unwrap())
             .file_io(FileIO::new_with_fs())
             .metadata_location("/test/metadata.json")
             .runtime(test_runtime())
             .build()
-            .unwrap();
-
-        let result = project_with_partition(input, &table);
-        assert!(result.is_ok(), "Schema validation should pass");
-    }
-
-    #[test]
-    fn test_schema_validation_mismatched_schemas() {
-        use iceberg::TableIdent;
-        use iceberg::io::FileIO;
-        use iceberg::spec::{FormatVersion, NestedField, PrimitiveType, Schema, 
Type};
-
-        let table_schema = Arc::new(
-            Schema::builder()
-                .with_fields(vec![
-                    NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int))
-                        .into(),
-                    NestedField::required(
-                        2,
-                        "name",
-                        Type::Primitive(PrimitiveType::String),
-                    )
-                    .into(),
-                ])
-                .build()
-                .unwrap(),
-        );
-
-        let partition_spec = PartitionSpec::builder(table_schema.clone())
-            .add_partition_field("id", "id_partition", Transform::Identity)
             .unwrap()
-            .build()
-            .unwrap();
+    }
 
-        let sort_order = iceberg::spec::SortOrder::builder()
-            .build(&table_schema)
-            .unwrap();
+    /// Required `id: int` and `name: string` columns.
+    fn id_and_name_table() -> Table {
+        table_partitioned_by_id(vec![
+            NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)),
+            NestedField::required(2, "name", 
Type::Primitive(PrimitiveType::String)),
+        ])
+    }
 
-        let table_metadata_builder = iceberg::spec::TableMetadataBuilder::new(
-            (*table_schema).clone(),
-            partition_spec,
-            sort_order,
-            "/test/table".to_string(),
-            FormatVersion::V2,
-            std::collections::HashMap::new(),
-        )
-        .unwrap();
+    fn input_of(fields: Vec<Field>) -> Arc<dyn ExecutionPlan> {
+        Arc::new(EmptyExec::new(Arc::new(ArrowSchema::new(fields))))
+    }
 
-        let table_metadata = table_metadata_builder.build().unwrap();
+    const INCOMPATIBLE: &str = "Input schema is not compatible with Iceberg 
table schema";
 
-        // Create Arrow schema with different field name (mismatched)
-        let arrow_schema = Arc::new(ArrowSchema::new(vec![
+    #[test]
+    fn test_schema_validation_matching_schemas() {
+        let input = input_of(vec![
             Field::new("id", DataType::Int32, false),
-            Field::new("different_name", DataType::Utf8, false), // Wrong 
field name
-        ]));
+            Field::new("name", DataType::Utf8, false),
+        ]);
+        assert!(project_with_partition(input, &id_and_name_table()).is_ok());
+    }
 
-        let input = Arc::new(EmptyExec::new(arrow_schema));
+    #[test]
+    fn test_schema_validation_mismatched_schemas() {
+        let input = input_of(vec![
+            Field::new("id", DataType::Int32, false),
+            Field::new("different_name", DataType::Utf8, false),
+        ]);
+        let err = project_with_partition(input, &id_and_name_table())
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
+    }
 
-        let table = Table::builder()
-            .metadata(table_metadata.metadata)
-            .identifier(TableIdent::from_strs(["test", "table"]).unwrap())
-            .file_io(FileIO::new_with_fs())
-            .metadata_location("/test/metadata.json")
-            .runtime(test_runtime())
-            .build()
-            .unwrap();
+    #[test]
+    fn test_schema_validation_nullability() {
+        let id = |nullable| input_of(vec![Field::new("id", DataType::Int32, 
nullable)]);
+        let int = Type::Primitive(PrimitiveType::Int);
+
+        // A non-nullable input fits an optional column, e.g. an INSERT from a
+        // NOT NULL source.
+        let optional =
+            table_partitioned_by_id(vec![NestedField::optional(1, "id", 
int.clone())]);
+        assert!(project_with_partition(id(false), &optional).is_ok());
+        assert!(project_with_partition(id(true), &optional).is_ok());
+
+        // A nullable input could write nulls into a required column.
+        let required = table_partitioned_by_id(vec![NestedField::required(1, 
"id", int)]);
+        assert!(project_with_partition(id(false), &required).is_ok());
+        let err = project_with_partition(id(true), &required)
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
+    }
 
-        let result = project_with_partition(input, &table);
-        assert!(
-            result.is_err(),
-            "Schema validation should fail for mismatched schemas"
-        );
-        assert!(
-            result
-                .unwrap_err()
-                .to_string()
-                .contains("Input schema does not match Iceberg table schema")
-        );
+    #[test]
+    fn test_schema_validation_nested_nullability() {
+        let child = |nullable| Field::new("x", DataType::Int32, nullable);
+        let input = |nullable| {
+            input_of(vec![
+                Field::new("id", DataType::Int32, false),
+                Field::new(
+                    "s",
+                    DataType::Struct(Fields::from(vec![child(nullable)])),
+                    false,
+                ),
+            ])
+        };
+        let table = |x: NestedField| {
+            table_partitioned_by_id(vec![
+                NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)),
+                NestedField::required(
+                    2,
+                    "s",
+                    Type::Struct(StructType::new(vec![Arc::new(x)])),
+                ),
+            ])
+        };
+        let int = Type::Primitive(PrimitiveType::Int);
+
+        // The same rule applies inside a struct.
+        let optional = table(NestedField::optional(3, "x", int.clone()));
+        assert!(project_with_partition(input(false), &optional).is_ok());
+        assert!(project_with_partition(input(true), &optional).is_ok());
+
+        let required = table(NestedField::required(3, "x", int));
+        assert!(project_with_partition(input(false), &required).is_ok());
+        let err = project_with_partition(input(true), &required)
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains(INCOMPATIBLE), "{err}");
     }

Review Comment:
   Filed apache/iceberg-rust#3297 for the `strip_metadata_from_schema` bug. The 
fix plus a regression test is ready locally, and I'll open the PR after 
apache/iceberg-rust#3286 and apache/iceberg-rust#2904 are in.
   
   The staged list and map tests are 
[here](https://github.com/apache/datafusion-iceberg/blob/abee0425e1c7af8ff0f5b323209c6d36e27fa4d1/crates/datafusion/src/physical_plan/project.rs#L831-L878).
 Once datafusion-iceberg is on an iceberg-rust rev with the fix, enabling them 
is just uncommenting, since they already pass against `665c64e` with the fix 
applied.
   
   Two list/map limitations of `contains` I noticed while writing them, neither 
new in this PR (the old `==` had both):
   - Field names are compared, so a list whose element is named `item` (Arrow's 
default, used by `DataType::new_list`) won't match Iceberg's `element`.
   - A source map with `keys_sorted = true` is rejected, since Iceberg maps 
convert to unsorted Arrow maps.
   
   If you think either is worth handling, I can fold them into #22 or open a 
separate issue.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to