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]

Reply via email to