fix: accept NOT NULL input columns in partitioned INSERT - #19
NoahKusaba wants to merge 4 commits into
Conversation
mbutrovich
left a comment
There was a problem hiding this comment.
Thanks @NoahKusaba! The fix looks right to me, and I confirmed that test_insert_not_null_source_into_partitioned_table fails with main's project.rs and passes with this change.
While testing the partitioned and unpartitioned paths side by side, I found a gap on the unpartitioned side that this PR doesn't cause but that sits next to its motivation. Could you open an issue for it? The unpartitioned path skips project_with_partition (table/mod.rs), and IcebergWriteExec passes the input's own schema as the sink schema to execute_input_stream (write.rs). DataFusion only runs its runtime null check (check_not_null_constraints) for columns that are non-nullable in the sink schema and nullable in the input. Because both schemas are the input's here, the check never runs. At the head commit, inserting SELECT * FROM source into an unpartitioned table with a required id: int column, from a MemTable whose nullable id holds [1, NULL], succeeds and reads back 0, 1. The NULL is written as 0. The partitioned path goes the other way and rejects a nullable source at plan time even when it holds no nulls, which DataFusion's own sinks accept and check at runtime. Passing the table's Arrow schema (plus the partition column) as the sink schema would probably fix both, so the issue could cover both.
| 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}"); | ||
| } |
There was a problem hiding this 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, and the map arm also compares the keys_sorted flag, so a struct test doesn't cover them.
There was a problem hiding this 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 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 byDataType::new_list) won't match Iceberg'selement. - A source map with
keys_sorted = trueis 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.
Move the four nullability cases into assert_nested_nullability and run them for a struct field. Add the same cases for a list element and a map value, commented out until iceberg-rust's strip_metadata_from_schema supports lists and maps; today it errors on them before the schemas are compared.
|
@mbutrovich Opened #22 for the unpartitioned gap, with your repro. It covers both sides:
I'd like to keep this PR to the narrower fix, accepting NOT NULL sources into optional columns, and leave the plan-time rejection as it is. Once #22 passes the table's schema as the sink schema, that rejection can be dropped in favour of DataFusion's runtime check, and both paths will behave the same. |
andygrove
left a comment
There was a problem hiding this comment.
Thanks @NoahKusaba! The contains swap looks right to me. With metadata stripped from both sides, nullability is the only thing it relaxes, and in the right direction: a nullable input into a required column is still rejected. The PR also merges cleanly onto current main (including #21's iceberg-rust bump), and fmt, clippy and the workspace tests pass there. I think two things need sorting out before merging, though. I reproduced both against main at a2bc942 and with this PR merged onto it.
1. This lets NOT NULL sources reach #35
In #35, when the insert input ends in a projection that reorders columns, projection pushdown merges iceberg_partition_values into that projection. PartitionExpr declares no children, so it then reads columns by position from the batch before the reorder. On main, the equality check here happens to reject NOT NULL sources into optional columns before they reach that. With this PR they get through and silently write the wrong partitions.
Table t has optional id and val and is partitioned by identity(id). Source src(a, b) holds (1, 100), (2, 200):
| Insert | main |
This PR |
|---|---|---|
INSERT INTO t SELECT b AS id, a AS val FROM src, NOT NULL src |
planning error | succeeds, but the data files land in id=1/ and id=2/, and SELECT count(*) FROM t WHERE id = 100 returns 0 |
Same, nullable src |
wrong partitions (#35) | wrong partitions (#35) |
INSERT INTO t (val, id) VALUES (1, 100), (2, 200) |
wrong partitions (#35) | wrong partitions (#35) |
INSERT INTO t VALUES (100, 1), (200, 2) |
correct | correct |
The physical plan for the (val, id) VALUES case shows the merge:
IcebergWriteExec: table=ns.t
ProjectionExec: expr=[column2@1 as id, column1@0 as val, iceberg_partition_values(id) as _partition]
DataSourceExec: partitions=1, partition_sizes=[1]
This is #35's bug rather than this PR's. Still, I think #35 should be fixed first, or in the same change, so this PR doesn't open a new path to silent corruption. Optional is the common case for Iceberg columns.
2. Nested nullability passes planning but fails in the Parquet writer
The new doc says a column may be non-nullable "at any nesting depth", but the write path only accepts top-level differences. iceberg-rust's ParquetWriter builds its Arrow writer from the table schema. parquet 59.3's LevelInfoBuilder::types_compatible compares types with equals_datatype, which includes nested nullability.
Take a struct column c whose child x is NOT NULL in the source and optional in the table, inserted through DataFrame::write_table:
main:Error during planning: Input schema does not match Iceberg table schema.- This PR:
Failed to write using parquet writer., source: Arrow: Incompatible type. Field 'c' has type Struct("x": Int32, metadata: {"PARQUET:field_id": "3"}), array has type Struct("x": non-null Int32)
SQL INSERT ... SELECT doesn't hit this, because DataFusion casts the column to the table's struct type. Unpartitioned tables already fail the same way at write time since they skip this check, so that part isn't new.
Relaxing only top-level nullability would keep the fix for #27 and still reject nested differences at planning:
let fits = input_schema_cleaned.fields().len() == expected_schema_cleaned.fields().len()
&& input_schema_cleaned
.fields()
.iter()
.zip(expected_schema_cleaned.fields())
.all(|(input, expected)| {
input.name() == expected.name()
&& input.data_type() == expected.data_type()
&& (expected.is_nullable() || !input.is_nullable())
});I tried that, with test_schema_validation_struct_nullability changed to expect an error for the NOT NULL → optional case:
- The unit and integration tests pass.
- The nested
write_tableinsert is rejected at planning again. - Top-level NOT NULL inserts work through both SQL and
write_table.
The staged list and map tests would need the same change. Happy to share the probe tests if they're useful.
| // 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 |
There was a problem hiding this comment.
Could these be #[ignore = "apache/iceberg-rust#3297"] tests instead of comments? They'd keep compiling, and enabling them would just mean removing the attribute.
There was a problem hiding this 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.
I'll set the PR to draft as well in the meantime. We can also wait for #35 to get merged in before too :)
|
@andygrove Thank you for taking the time out of your busy schedule to deeply review my change, and for supporting the ballista-iceberg initiative. I need to spend more time digesting your comment, before I can give a more proper response (if I need to make one). Also your book is amazing thank you for publishing it free online <3 |
The closure inside the helper is also named input and shadowed the parameter it calls.
These depend on iceberg-rust's strip_metadata_from_schema handling list and map columns (apache/iceberg-rust#3303) and fail until that lands.
Which issue does this PR close?
What changes are included in this PR?
An
INSERTinto a partitioned table fails when the source has aNOT NULLcolumn where the table's column is optional:Every value of a non-nullable column is valid in an optional one, so the write is safe. The same
INSERTinto an unpartitioned table already succeeds, becauseproject_with_partitionreturns before this check when the spec is unpartitioned, so today the result depends on whether the table is partitioned.project_with_partitioncompared the input and table schemas with==. It now uses Arrow'sSchema::contains, which allows the input to be narrower in nullability and is otherwise as strict as before:The public doc of
project_with_partitionnow states this contract. The error message changes from "does not match" to "is not compatible with", since the schemas no longer need to be equal.The new unit tests need a partitioned table with a given schema. The three existing schema-validation tests each built one inline with the same ~50 lines, so they now share a
table_partitioned_by_idhelper, with their imports moved to the top of thetestsmodule. Their assertions are unchanged apart from the new error message.Are these changes tested?
test_insert_not_null_source_into_partitioned_table(integration): inserts from aNOT NULLMemTableinto a table partitioned on an optional column, and reads the rows back. It fails onmainwith the error above.test_schema_validation_nullability: non-nullable and nullable input into an optional column are both accepted; into a required column, non-nullable is accepted and nullable is rejected.test_schema_validation_struct_nullability: the same four cases for a field inside a struct.cargo fmt --all -- --check,cargo clippy --workspace --locked --all-targets -- -D warningsandcargo test --workspace --lockedall pass locally.AI Disclosure
main.🤖 Generated with Claude Code