Neuw84 commented on issue #6264: URL: https://github.com/apache/datafusion-comet/issues/6264#issuecomment-5884461591
A related symptom, possibly the same root cause: with the native scan, a dynamic partition pruning filter can be dropped from a scan, so the scan reads every partition. Results are still correct; the cost is the extra reading. I'm reporting it here first because it sits in the same area. It may deserve its own issue; I can open one. **Where it shows up.** TPC-DS q23a/q23b at SF1000, Parquet with the fact tables partitioned by `*_sold_date_sk`. Inside the `max_store_sales` scalar subquery, the `store_sales` scan carries `dynamicpruningexpression(ss_sold_date_sk IN dynamicpruning#…)` from the `d_year IN (2000, 2001, 2002, 2003)` date scan. | that `store_sales` scan, SF1000 | Spark's `FileSourceScanExec` | `CometNativeScanExec` (Comet 1.0.0) | |---|---|---| | DPP subquery attached to the scan in the final plan | `SubqueryBroadcast dynamicpruning#…` | none | | rows / file ranges read | 1.66 B rows (pruned) | 2.75 B rows, 14,584 of 14,616 file ranges (unpruned) | | bytes scanned | | 9.7 GB | The other DPP scans in the same query (`frequent_ss_items`, and the `catalog_sales` and `web_sales` scans) keep their subquery under Comet and prune like Spark. **When it happens.** It depends on the join shape: - The failing scan sits under a shuffle for a shuffled join with `customer`, and the broadcast join with `date_dim` that the DPP filter reuses is above that join. - At SF1 with default settings `customer` is broadcast too, and the native scan prunes correctly. **Minimal reproduction (SF1).** 1. TPC-DS SF1 Parquet with `store_sales`, `catalog_sales` and `web_sales` written `partitionBy(<x>_sold_date_sk)`; dimensions unpartitioned. 2. Settings: `spark.sql.autoBroadcastJoinThreshold=200000`, so `customer` is shuffled and `date_dim` stays broadcast. AQE is on (default). `spark.comet.scan.enabled=true` with `spark.comet.scan.impl` left as the default (the native scan). 3. Run q23a. | SF1, threshold 200 KB | the scan's DPP subquery | file ranges read (of 14,592) | checksum | |---|---|---|---| | Spark only | `SubqueryBroadcast` | pruned (1.66 M rows) | `36aec83f` | | Comet scan, Spark operators | none | 14,584 | `36aec83f` | | Comet scan + Comet exec and shuffle (`CometSortMergeJoin`, `CometBroadcastHashJoin`, `CometExchange`) | none | 14,584 | `36aec83f` | The 8 file ranges that are skipped are the `__HIVE_DEFAULT_PARTITION__` ones, removed by the static `isnotnull(ss_sold_date_sk)`. The DPP filter itself prunes nothing. **Relation to this issue.** I could not tie it directly to `filterUnusedDynamicPruningExpressions`. In 1.0.0 that function is only called from `CometNativeScanExec.doCanonicalize` (and the scan here is not reused; it runs with its own metrics). What the plan shows is the DPP subquery missing from the native scan, so the scan's partition list never gets the runtime filter. If an unconverted `SubqueryAdaptiveBroadcastExec` is being dropped here outside canonicalization as well, the fix for this issue may cover it; if not, a separate issue is probably right. -- 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]
