andygrove opened a new pull request, #6543:
URL: https://github.com/apache/datafusion-comet/pull/6543

   ## Which issue does this PR close?
   
   Closes #6504.
   
   ## Rationale for this change
   
   The native Iceberg scan reads files through iceberg-rust, which does not 
match a data file's nested fields to the table's by field id 
(apache/iceberg-rust#2617, with a fix in review in apache/iceberg-rust#3255). 
Two kinds of nested schema evolution break it.
   
   A nested field added after a file was written, to a struct or to a struct 
inside an array or map, fails the scan with `Incorrect number of arrays for 
StructArray fields`. iceberg-rust casts the file's struct to the table's wider 
struct, and the cast cannot add a field. That is #6504. A plain `SELECT s` 
already failed this way in 1.0.0. In 1.1.0, #5732 removed the fallback for `IS 
NULL` and `IS NOT NULL` on complex columns, so null checks, non-outer 
`explode`, `posexplode` and `inline`, and joins on such a column now fail where 
1.0.0 fell back and returned the right rows. Spark's nested pruning does not 
help, because the native read asks for the column's full nested type, so 
`SELECT s.a` fails too.
   
   A nested field renamed after a file was written can read back as NULL, with 
no error. When the file's nested fields have the same nullability as the 
table's, iceberg-rust passes the column through with the file's old field 
names. Comet's batch adaptation then looks struct fields up by name and fills 
the renamed one with NULL. A reorder combined with a rename returns wrong 
values in either layout. This bug is separate from #5732 and has no issue of 
its own. It turned up while writing the tests: Spark 4 writes a struct built 
only from literals with `required` fields, which takes iceberg-rust's cast path 
and hides the bug, while Spark 3.4, or any row with a null in the struct, 
writes `optional` fields and shows it.
   
   ## What changes are included in this PR?
   
   `CometScanRule` declines the native scan when a projected column has a 
nested field that some schema in the table's history lacks or names 
differently. A `FileScanTask` does not record which schema wrote its file, so 
the schema history is the evidence. The new 
`IcebergReflection.nestedFieldsAddedOrRenamed` matches fields by id at every 
level (struct fields, array elements, map keys and values). It checks the two 
schemas the native read takes a column's type from: the current table schema, 
and the scan schema that is used when `VERSION AS OF` reads a dropped column. 
The fallback reason names each field, for example `s.b (added)` or 
`items.element.z (renamed from a)`.
   
   Nested drops, reorders without a rename, and type promotions still read 
natively, because they read correctly today. The read projects leaves by field 
id, so a dropped field is never read, and reorders and promotions keep the 
field names, which the reader resolves correctly. Scans that don't project an 
evolved column are unaffected.
   
   The check is conservative. Old schemas stay in the history, so a table keeps 
falling back for an evolved column after compaction rewrites the old files, and 
also when the change was made before any data was written. The check can go 
once Comet's iceberg-rust pin includes apache/iceberg-rust#3255.
   
   The projected field ids were computed inside the Variant check and are now 
hoisted so both checks share them. The Iceberg user guide lists the new 
fallback under current limitations.
   
   ## How are these changes tested?
   
   There are five new tests in `CometIcebergNativeSuite`. Each fallback test 
compares with Spark, checks the fallback reason, and asserts that no native 
Iceberg scan is planned.
   
   - A field added to a struct: `SELECT s`, `IS NOT NULL`, `IS NULL`, `SELECT 
s.a`, a self-join on `s`, and `VERSION AS OF` a snapshot from before the 
change. `SELECT id` stays native.
   - A field added inside an array (`explode`, `posexplode`, `inline`), inside 
a map value, several levels down, and a struct-typed field.
   - A nested field dropped and re-added under the same name in one schema 
change, so that only the field id tells the two apart.
   - Renames in a struct and in an array element, and a reorder plus a rename. 
The struct holds nulls, so the files take the pass-through path.
   - Nested drops, reorders and promotions stay native and match Spark, over 
files in both nullability layouts.
   
   Before the fix I reproduced every query in the issue failing on `main`. With 
the guard disabled, all four fallback tests fail: three with the `Incorrect 
number of arrays` error, and the rename test with NULLs. With the fields 
matched by name instead of by id, the re-add test fails.
   
   Local runs:
   
   - The new tests pass on Spark 3.4, 3.5, 4.0 and 4.1.
   - On the default Spark 4.1 profile, `CometIcebergNativeSuite`, 
`CometFuzzIcebergSuite`, `CometIcebergRewriteActionSuite`, 
`CometIcebergResidualPushdownSuite`, `CometIcebergWriteActionSuite`, 
`CometIcebergWriteDetectionSuite`, `CometIcebergSystemFunctionSuite`, 
`IcebergReflectionSuite` and `CometIcebergNativeScanSuite` ran 318 tests that 
passed and 1 that was canceled, the existing SPARK-55626 case.
   
   I have not run the Iceberg Spark SQL tests locally, so this PR has the 
`run-iceberg-tests` label.
   


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