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]
