andygrove opened a new issue, #4988:
URL: https://github.com/apache/datafusion-comet/issues/4988

   ### What is the problem the feature request solves?
   
   Follow-up to #3457 with measurements and root-cause analysis.
   
   `spark.comet.parquet.rowFilterPushdown.enabled` (DataFusion's 
`pushdown_filters`, i.e. row-level `RowFilter` evaluation plus late 
materialization) still defaults to `false`. #4722 already made format-level 
pruning (row group statistics, page index, bloom filters) work by default, so 
this flag is now the only remaining piece of "filter pushdown" that is opt-in.
   
   I benchmarked TPC-H Q1 and Q6 at SF100 to see whether it can be turned on by 
default. It cannot, and the reason is structural rather than a tuning problem.
   
   ### Measurements
   
   Setup: TPC-H SF100 local Parquet (`lineitem` = 600,037,902 rows, 25 GB, 32 
files, 15,744 row groups), Spark 3.5.8 standalone, 2 executors x 8 cores, warm 
page cache, 10 iterations per config. Each config was run in both orders 
(pushdown off first, then on first) to control for cache-warming effects. 
Median of the last 5 iterations:
   
   | Query | off -> on (off ran first) | off -> on (on ran first) |
   | --- | --- | --- |
   | Q1 | 4.074 -> 3.726 s (-8.5%) | 3.804 -> 3.711 s (-2.4%) |
   | Q6 | 0.843 -> 1.523 s (+80.7%) | 1.029 -> 1.515 s (+47.2%) |
   
   Result hashes matched across all runs. The pushdown-on timings are stable in 
both orders (Q1 3.70-3.73 s, Q6 1.51-1.52 s); the off timings are what shift 
with cache state. Q1 is neutral to slightly positive, Q6 is a large 
reproducible regression.
   
   ### Root cause
   
   Scan metrics from the Spark event log, per iteration (600,037,902 rows):
   
   **Q6**
   
   | Metric | pushdown off | pushdown on |
   | --- | --- | --- |
   | Number of bytes scanned | 5,971,436,225 | 5,971,436,225 (identical) |
   | Scan output rows | 600,037,902 | 11,422,456 |
   | Data decompression + decoding (summed over tasks) | 7.89 s | 20.02 s |
   | Time evaluating row-level pushdown filters | ~0 | 0.37 s |
   | Row groups pruned by statistics | 0 of 15,744 | 0 of 15,744 |
   
   The filter itself works correctly: the scan emits 11.4M rows instead of 
600M, the expected ~1.9% selectivity. But **bytes scanned is byte-for-byte 
identical** and decode time is 2.5x higher.
   
   Late materialization can only avoid work when surviving rows are 
*clustered*. If every data page contains at least one surviving row, the reader 
still fetches and decompresses every page and skips only the final value-decode 
step. Q6's predicate (`l_shipdate` range AND `l_discount` range AND `l_quantity 
< 24`) is uncorrelated with physical row order, so ~1.9% selectivity is spread 
uniformly and no page is ever eliminated.
   
   What Q6 pays with the flag on:
   
   1. Full decode of the three filter columns (`l_shipdate`, `l_discount`, 
`l_quantity`) to evaluate the predicate. The predicate-cache metric confirms 
this exactly: 1,800,113,706 rows physically decoded per iteration = 3 filter 
columns x 600,037,902 rows.
   2. Selection-mask construction and predicate-cache bookkeeping.
   3. Selection-based decode of `l_extendedprice`, the only column late 
materialization can defer.
   4. Re-evaluation of the same predicate in the `CometFilter` above the scan 
(see item 1 in "Possible work items").
   
   against a baseline that simply decodes four columns straight through. That 
is strictly more work. The predicate cache is functioning (14.4B cached reads 
across 3 iterations), it just has nothing to skip.
   
   Q6 has the worst possible shape for this feature: the filter columns *are* 
three of the four projected columns, so the theoretical upside is deferring one 
column, while the costs apply to all of them.
   
   **Q1** is the mild version of the same effect. Its predicate (`l_shipdate <= 
date '1998-12-01' - interval '68 days'`) removes only 4.8M of 600M rows (0.8%), 
and decode goes 12.36 s -> 14.56 s summed over tasks. It comes out roughly 
neutral on wall clock.
   
   ### Side observation
   
   Statistics pruning eliminated **0 of 15,744 row groups** and the page index 
filtered 0 rows, in every configuration, for both queries. The format-level 
pruning that #4722 enabled by default does nothing on this dataset for these 
queries, for the same clustering reason. This is a property of unsorted TPC-H 
data rather than a Comet defect, but it means neither Q1 nor Q6 can demonstrate 
the upside of #4722 either, and it is worth keeping in mind when using TPC-H to 
evaluate any pruning work.
   
   ### Describe the potential solution
   
   The flag should stay opt-in. Enabling it globally regresses 
uncorrelated-predicate scans on un-clustered fact tables, which is the common 
case.
   
   Possible work items, roughly in order of value:
   
   1. **Drop the redundant `CometFilter` when every filter pushes down.** 
`native/core/src/parquet/parquet_exec.rs` deliberately discards 
`propagation.parent_pushdown_result` because Spark's `Filter` above the scan 
re-evaluates every data filter. With the flag on, the predicate is therefore 
evaluated twice. Plumbing the per-filter `PushedDown::Yes`/`No` classification 
back to the JVM would let `CometExecRule` remove the redundant filter. This 
reduces the overhead but does not by itself turn Q6 into a win.
   
   2. **Decouple `reorder_filters` from the flag.** `jni_api.rs` sets 
`pushdown_filters` and `reorder_filters` together. Upstream defaults 
`reorder_filters` to `false`; heuristic reordering has its own per-batch cost 
and can choose a worse order. These should be separately configurable so a 
regression can be attributed to the right one.
   
   3. **Expose `max_predicate_cache_size`.** It is currently never set, so 
DataFusion uses `None` (the arrow-rs default). The cache is what prevents 
filter columns being decoded twice and is worth being able to tune and measure.
   
   4. **Make this a planner decision rather than a global flag.** The feature 
wins when the projection is much wider than the filter columns and the 
predicate is selective, and loses when the filter columns dominate the 
projection or selectivity is uncorrelated with row order. A heuristic on 
projection width beyond the filter columns plus estimated selectivity would 
capture the upside without the Q6-shaped regressions.
   
   ### Additional context
   
   Measured on `apache/main` plus two unrelated aggregate-serde commits (no 
scan-path changes), DataFusion 54.0.0, arrow/parquet 58.3.0.
   
   Benchmarks were run with `benchmarks/tpc/run.py` using a `comet-rowfilter` 
engine variant that sets `spark.comet.parquet.rowFilterPushdown.enabled=true`.
   


-- 
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