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]

Reply via email to