pingzh opened a new pull request, #6263:
URL: https://github.com/apache/datafusion-comet/pull/6263

   ## Which issue does this PR close?
   
   Part 4 of #5775, following merged #6181 and the [five-PR 
plan](https://github.com/apache/datafusion-comet/issues/5775#issuecomment-5672549320).
 Leaves #5775 open for the separate decoded-batch filtering work in PR 5.
   
   ## Rationale for this change
   
   Local TopK fusion lets the heap and native Parquet scan run together. This 
PR connects the heap's improving threshold to the reader, allowing it to skip 
row groups that cannot improve the local candidates.
   
   ## What changes are included in this PR?
   
   Three commits:
   
   1. Connect an execution-local TopK predicate to supported native Parquet 
readers. Add the opt-in `spark.comet.exec.topK.dynamicFilter.enabled` setting 
and attachment counters. It requires the existing fusion option; both default 
to false.
   2. Cover lifecycle, cancellation, fallback, integer boundaries, 
dictionaries, nested timestamp errors, and Spark reader configurations.
   3. Extend the benchmark to isolate reader pruning from fusion and document 
eligibility, limitations, and metrics.
   
   The supported shape remains one direct signed integer sort key on the fused 
native Parquet scan. Each execution gets a fresh heap and predicate, including 
repeated, concurrent, reset, and replaced-child executions. Thresholds do not 
cross exchanges or Spark partitions. There is no added residual filter over 
decoded scan batches.
   
   Reader attachment preserves the existing per-file conversion guard. Fetch 
limits, supplied file statistics, and static predicates other than direct 
column null checks keep ordinary reading. Incomplete null counts remain 
unknown: a regression reproduced a wrong result before the safeguard (`0` 
instead of `NULL` for `ASC NULLS FIRST LIMIT 1`). The safeguard clones 
incomplete metadata without changing cached originals or page indexes. Join 
filtering is unchanged.
   
   Attachment counters belong to local TopK. Scan bytes, emitted rows, 
page/decoder counters, and the newly exposed `row_groups_pruned_dynamic_filter` 
remain on the scan.
   
   ## How are these changes tested?
   
   - Native runtime-filter tests: **50 passed**, including existing shared 
reader and join tests.
   - Spark 4.1 `CometTopKSuite`: **55 passed**, including filtering on/off, 
page-index and decoder-filter combinations, AQE, offsets, empty buckets, and 
structured corrupt-file paths.
   - Native release build and full Maven reactor install passed. Packaged JNI 
hash matches the freshly built library.
   - Scala formatting/style, Rust formatting, Markdown formatting, CI 
configuration validation, and `git diff --check` passed.
   - Expanded Spark SQL CI is pending. This draft needs a maintainer to apply 
`run-spark-4.1-tests` (and `run-benchmark-check` for the benchmark) before it 
is ready to merge; the author account has no label permissions.
   - Benchmark results below check full results against Spark, native plan 
shapes, actual partition counts, attachment counts, and reader metrics outside 
timed runs.
   
   ## Results
   
   **TL;DR:** Across this matrix, reader pruning made ascending layouts **3–65% 
faster than fused TopK without filtering**, with **46–82% fewer scan bytes**. 
Descending layouts pruned nothing; random layouts saved 0–10% of scan bytes 
with mixed timings. The existing wide-row descending fusion slowdown remains, 
so both options stay disabled by default.
   
   ### Setup
   
   - Spark 4.1, JDK 21, Linux, AMD EPYC 9V74 host (32 logical CPUs), Spark 
`local[1]`, 8 GiB JVM heap, native **release** build.
   - 1,048,576 rows; 1 or 4 scan partitions; 1 or 16 BIGINT payload columns; 
K=16 or 100,000; ascending, descending, and deterministic random layouts.
   - Snappy, dictionaries disabled, 5–6 verified row groups per file, batch 
size 4096, AQE off. Page-index and decoder filtering are off to isolate 
row-group pruning.
   - Two fresh JVM runs with opposite mode order. Each case warms up and 
measures for at least 500 ms and records at least 5 measured iterations. All 
**144 mode cases** completed; full Spark result, native-plan, partition, 
attachment, and metric checks passed.
   
   `unfused` disables fusion and filtering; `fused` enables only fusion; 
`pruning` enables both. Times below average the two runs' reported mean 
milliseconds and include query planning and collection. Scan counters come from 
untimed validation, matched across both runs, and use `bytes_scanned` 
(requested data/Bloom-filter ranges, excluding footer/page-index reads). There 
are no Bloom-filter reads in these fixtures.
   
   The larger K and multi-partition cases show why fewer bytes need not mean a 
proportional speedup: local heaps, the final TopK, and collecting the result 
still take time. For descending wide rows at K=16, fusion alone remains 29–36% 
slower than unfused execution. This PR does not remove that separate pipeline 
limitation.
   
   <details>
   <summary>All 24 scenarios (negative time change means faster)</summary>
   
   | Layout | P | Payload | K | Unfused ms | Fused ms | Pruning ms | Time 
change vs fused | Scan-byte saving |
   |---|---:|---:|---:|---:|---:|---:|---:|---:|
   | ascending | 1 | 1 | 16 | 68 | 54.5 | 42.5 | -22.0% | 81.6% |
   | ascending | 1 | 1 | 100000 | 107 | 105 | 95.5 | -9.0% | 81.6% |
   | descending | 1 | 1 | 16 | 165 | 163.5 | 168.5 | +3.1% | 0.0% |
   | descending | 1 | 1 | 100000 | 429 | 416.5 | 418 | +0.4% | 0.0% |
   | random | 1 | 1 | 16 | 48.5 | 43 | 46.5 | +8.1% | 8.2% |
   | random | 1 | 1 | 100000 | 582 | 554 | 564.5 | +1.9% | 0.0% |
   | ascending | 1 | 16 | 16 | 220 | 214.5 | 75.5 | -64.8% | 77.6% |
   | ascending | 1 | 16 | 100000 | 360 | 361 | 202 | -44.0% | 77.6% |
   | descending | 1 | 16 | 16 | 300 | 408.5 | 407 | -0.4% | 0.0% |
   | descending | 1 | 16 | 100000 | 655 | 756.5 | 761 | +0.6% | 0.0% |
   | random | 1 | 16 | 16 | 228 | 212 | 201.5 | -5.0% | 10.3% |
   | random | 1 | 16 | 100000 | 897.5 | 968.5 | 969.5 | +0.1% | 0.0% |
   | ascending | 4 | 1 | 16 | 70 | 59.5 | 47.5 | -20.2% | 82.1% |
   | ascending | 4 | 1 | 100000 | 314.5 | 325 | 315 | -3.1% | 46.3% |
   | descending | 4 | 1 | 16 | 189.5 | 183 | 186.5 | +1.9% | 0.0% |
   | descending | 4 | 1 | 100000 | 577.5 | 573 | 558 | -2.6% | 0.0% |
   | random | 4 | 1 | 16 | 72 | 55.5 | 62 | +11.7% | 5.2% |
   | random | 4 | 1 | 100000 | 795.5 | 801 | 794 | -0.9% | 0.0% |
   | ascending | 4 | 16 | 16 | 270.5 | 244.5 | 90 | -63.2% | 80.3% |
   | ascending | 4 | 16 | 100000 | 781 | 813.5 | 703.5 | -13.5% | 60.7% |
   | descending | 4 | 16 | 16 | 347 | 446.5 | 447.5 | +0.2% | 0.0% |
   | descending | 4 | 16 | 100000 | 1046.5 | 1121 | 1114.5 | -0.6% | 0.0% |
   | random | 4 | 16 | 16 | 241 | 233.5 | 233.5 | +0.0% | 1.2% |
   | random | 4 | 16 | 100000 | 1313 | 1424 | 1436 | +0.8% | 0.0% |
   
   </details>
   
   Reproduce with a release build:
   
   ```sh
   make PROFILES=-Pspark-4.1 BENCH_HEAP=8g 
benchmark-org.apache.spark.sql.benchmark.CometTopKBenchmark -- 1048576 1,4 1,16 
16,100000 ascending,descending,random 5 unfused,fused,pruning
   ```
   
   Repeat with final argument `pruning,fused,unfused`. The local runs used the 
equivalent cached Maven benchmark command after a full reactor install. Both 
exited 0 and printed `BUILD SUCCESS`; Maven's `exec:java` runner emitted a 
Hadoop shutdown-hook classloader warning afterward.
   
   


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