mbutrovich commented on code in PR #19:
URL: https://github.com/apache/datafusion-iceberg/pull/19#discussion_r4126428819
##########
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:
The new doc says the rule holds "at any nesting depth", and this test covers
a struct field. Could we add the same four cases for a list element and a map
value? Lists and maps go through their own arms of Arrow's
[`DataType::contains`](https://github.com/apache/arrow-rs/blob/f90e061326bd821a7af09281d9e92de6f3b603d9/arrow-schema/src/datatype.rs#L817-L838),
and the map arm also compares the `keys_sorted` flag, so a struct test doesn't
cover them.
--
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]