comphead opened a new pull request, #6725: URL: https://github.com/apache/datafusion-comet/pull/6725
## Which issue does this PR close? Closes #6724. ## Rationale for this change The native Iceberg scan reads and decodes every field of a projected struct, list, or map column, even when Spark's nested schema pruning keeps only some of them. On TPC-H SF1000 with a nested schema, Q20's `partsupp` scan reads 26.4 GiB where the query needs only `ps_availqty`. This is the main reason Comet's scan lead shrinks on that schema. #6724 has the measurements. ## What changes are included in this PR? - `CometIcebergNativeScan` passes iceberg-rust the pruned scan schema (`SparkScan.expectedSchema()`) as the task schema, without its metadata columns, instead of the table schema. iceberg-rust then reads only the leaves the query uses. Tasks with deletes use it too. - Partition source columns and equality-delete keys are still added through `schemaWithRequiredFields`. The result is memoized by base schema and field ids, so tasks that need the same columns share one schema in the pool instead of serializing one each. - A task keeps the full schema when a missing partition source or equality-delete key is nested inside a struct the query pruned. Appending that field at the top level could clash with a top-level column of the same name (`Invalid schema: multiple fields for name region`). - A new config, `spark.comet.scan.icebergNative.nestedSchemaPruning.enabled` (default `true`), turns the pruning off. - `IcebergReflection` gains `withoutMetadataColumns` and `isNestedField`. - The nested-evolution fallback from #6543 is unchanged. Its comment in `CometScanRule`, and one in `CometIcebergNativeSuite`, now describe the pruned read. - `CometIcebergReadBenchmark` gets a nested column case. `--nested-only` runs only that case. ## How are these changes tested? Four new tests in `CometIcebergNativeSuite`: - Struct, nested struct, list, and map value pruning, NULL validity, and a nested filter, all checked against Spark. Also a `bytes_scanned` comparison with the config on and off. - Equality and position deletes on keys the query does not project, a partition source it does not project, and `VERSION AS OF`. Also a nested partition source next to a top-level column of the same name, which fails without the guard. - One task schema in the pool across eight partitions. - Data files without field ids and no name mapping. Spark's own reader returns NULL nested fields there, so the pruned native read is compared with the full native read. These ran locally on the default Spark 4.1 profile with a release native build: - the four new tests - `CometIcebergNativeSuite`: 120 passed, 1 version-gated cancel - `CometFuzzIcebergSuite`, `CometIcebergResidualPushdownSuite`, and `CometIcebergRewriteActionSuite` Not run: the S3 suites, Spark's SQL tests, and the other Spark profiles. `CometIcebergReadBenchmark --nested-only` reads 4M rows on one local core with a warm page cache: | Query | Spark | Comet, every nested field | Comet, pruned | |---|---:|---:|---:| | `explode(partsupp_data.ps_availqty)` (Q20 shape) | 413 ms | 624 ms | 62 ms | | `explode(partsupp_data.ps_comment)` (control) | 2,800 ms | 590 ms | 542 ms | Bytes scanned for the first query drop from 124.0 MiB to 7.5 MiB. TPC-H Q14, Q18, and Q20 on the nested schema at scale are still to be run. -- 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]
