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


##########
crates/datafusion/src/physical_plan/project.rs:
##########
@@ -675,215 +688,203 @@ 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
-        ]));
-
-        let input = Arc::new(EmptyExec::new(arrow_schema));
-
-        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();
-
-        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")
-        );
+            Field::new("name", DataType::Utf8, false),
+        ]);
+        assert!(project_with_partition(input, &id_and_name_table()).is_ok());
     }
 
     #[test]
-    fn test_schema_validation_with_metadata_differences() {
-        use std::collections::HashMap;
-
-        use iceberg::TableIdent;
-        use iceberg::io::FileIO;
-        use iceberg::spec::{FormatVersion, NestedField, PrimitiveType, Schema, 
Type};
+    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_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(),
-        );
+    #[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 partition_spec = PartitionSpec::builder(table_schema.clone())
-            .add_partition_field("id", "id_partition", Transform::Identity)
-            .unwrap()
-            .build()
-            .unwrap();
+    /// Checks the nullability rule for a nested field: `input(nullable)` has a
+    /// required `id` and a column `c` whose nested int is `nullable`, and
+    /// `nested_table(required)` has the matching table column.
+    fn assert_nested_nullability(
+        input: impl Fn(bool) -> DataType,
+        nested_table: impl Fn(bool) -> Type,
+    ) {
+        let input = |nullable| {
+            input_of(vec![
+                Field::new("id", DataType::Int32, false),
+                Field::new("c", input(nullable), false),
+            ])
+        };
+        let table = |required| {
+            table_partitioned_by_id(vec![
+                NestedField::required(1, "id", 
Type::Primitive(PrimitiveType::Int)),
+                NestedField::required(2, "c", nested_table(required)),
+            ])
+        };
 
-        let sort_order = iceberg::spec::SortOrder::builder()
-            .build(&table_schema)
-            .unwrap();
+        let optional = table(false);
+        assert!(project_with_partition(input(false), &optional).is_ok());
+        assert!(project_with_partition(input(true), &optional).is_ok());
 
-        let table_metadata_builder = iceberg::spec::TableMetadataBuilder::new(
-            (*table_schema).clone(),
-            partition_spec,
-            sort_order,
-            "/test/table".to_string(),
-            FormatVersion::V2,
-            HashMap::new(),
-        )
-        .unwrap();
+        let required = table(true);
+        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}");
+    }
 
-        let table_metadata = table_metadata_builder.build().unwrap();
+    #[test]
+    fn test_schema_validation_struct_nullability() {
+        let int = Type::Primitive(PrimitiveType::Int);
+        assert_nested_nullability(
+            |nullable| {
+                DataType::Struct(Fields::from(vec![Field::new(
+                    "x",
+                    DataType::Int32,
+                    nullable,
+                )]))
+            },
+            |required| {
+                let x = NestedField::new(3, "x", int.clone(), required);
+                Type::Struct(StructType::new(vec![Arc::new(x)]))
+            },
+        );
+    }
 
-        // Create Arrow schema with metadata (should be ignored in comparison)
-        let mut metadata = HashMap::new();
-        metadata.insert("extra".to_string(), "metadata".to_string());
+    // TODO: enable once iceberg-rust's `strip_metadata_from_schema` handles 
lists

Review Comment:
   I actually opened the PR on Iceberg-rust to let us enable these tests. 
   
   It's getting better traction than I thought thanks to compheads help, so I 
think we should just merge these tests in enabled when the iceberg-rust changes 
are in.
   
   https://github.com/apache/iceberg-rust/pull/3303
   
   I'll set the PR to draft as well in the meantime. We can also wait for #35  
to get merged in before too :) 



-- 
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