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]
