pingzh commented on issue #5775:
URL: 
https://github.com/apache/datafusion-comet/issues/5775#issuecomment-5920801524

   ## Completion and benchmark results
   
   **TL;DR:** On a fresh release build of merged `main`, ascending layouts 
reduced scan data bytes by **46–82%**. For one partition, 16 BIGINT payload 
columns, and K=16, pruning reduced mean query time from **214.5 ms to 76 ms 
versus fused TopK** (64.6% lower), and from **230.5 ms to 76 ms versus unfused 
TopK** (67.0% lower). Descending layouts pruned nothing; the wide K=16 cases 
remained **about 31% slower than unfused execution** because the fusion cost 
remains. Both options stay disabled by default.
   
   The [updated four-PR 
plan](https://github.com/apache/datafusion-comet/issues/5775#issuecomment-5672549320)
 is complete: #5937, #6067, #6181, and #6263 are merged. The former optional 
PR5 is no longer part of this issue's scope; #6263's older description still 
refers to that earlier plan.
   
   ### Setup
   
   - Revision: 
[`9c7fcc5aa4a343f53b4bfc8707eeb8dfd17a2c48`](https://github.com/apache/datafusion-comet/commit/9c7fcc5aa4a343f53b4bfc8707eeb8dfd17a2c48),
 freshly fetched Apache `main` on September 30, 2026; includes the merged 
reader-pruning implementation.
   - Spark **4.1.3**, Temurin **21.0.12.1**, DataFusion **55.1.0**, 
Arrow/Parquet Rust **59.3.0**; Linux, AMD EPYC 9V74 host with 32 logical CPUs. 
Native **release** build with `target-cpu=native`, 8 GiB JVM heap.
   - **1,048,576 rows**, one INT key, **1 or 16 BIGINT payload columns**, **1 
or 4 scan partitions**, **K=16 or 100,000**.
   - **Ascending, descending, and deterministic random physical layouts. Every 
query uses `SELECT * FROM topk_input ORDER BY k ASC LIMIT K`.** These layout 
names do not describe different SQL sort directions.
   - Spark `local[1]`, AQE off, one Parquet file per scan partition, **5–6 
verified row groups per file**, Snappy, dictionaries disabled, batch size 4096. 
Four partitions exercise separate task thresholds, not parallel scaling.
   - Page-index pruning and decoder row filtering disabled to isolate row-group 
pruning.
   - **Two fresh JVM runs with opposite mode order**: `unfused,fused,pruning`, 
then `pruning,fused,unfused`. Each case has at least 500 ms warmup, at least 
500 ms measurement, and at least five measured iterations: **24 scenarios × 3 
modes × 2 runs = 144 mode cases**.
   
   `unfused` disables both options; `fused` enables only 
`spark.comet.exec.topK.fusion.enabled`; `pruning` enables fusion plus 
`spark.comet.exec.topK.dynamicFilter.enabled`.
   
   Times below average the two runs' reported mean milliseconds and include 
query planning and collection. Fixture generation, result comparisons, and 
plan/metric checks are outside timings. These are warmed local microbenchmarks; 
small timing differences should not be treated as established regressions or 
speedups.
   
   ### Representative results
   
   One partition, 16 BIGINT payload columns, K=16:
   
   | Physical layout | Unfused ms | Fused ms | Pruning ms | Scan data bytes, 
off → on | Scan output rows, off → on | Row groups skipped |
   |---|---:|---:|---:|---:|---:|---:|
   | ascending | 230.5 | 214.5 | 76 | 73,485,975 → 16,482,556 | 1,048,576 → 
235,199 | 4 / 5 |
   | descending | 306.5 | 403 | 401 | 73,485,972 → 73,485,972 | 1,048,576 → 
1,048,576 | 0 / 5 |
   | random | 215.5 | 201 | 182 | 73,486,406 → 65,932,698 | 1,048,576 → 940,796 
| 1 / 5 |
   
   Across the full matrix:
   
   - Ascending layouts saved **46.3–82.1% of scan data bytes**. The wide K=16 
cases were **63.2–64.6% faster than fused without pruning**. Savings did not 
always reduce elapsed time: with four partitions, one payload column, and 
K=100,000, bytes fell **46.3%** while time was effectively unchanged (**313.5 → 
315.5 ms**).
   - Descending layouts saved **no scan bytes or row groups**. The wide K=16 
cases with fusion and pruning remained **30.8–31.2% slower than unfused**. 
Turning reader filtering on/off within the fused plan changed much less; the 
larger cost is already present with fusion alone.
   - Random layouts saved **0–10.3% of scan bytes**, with mixed timing changes. 
No bytes were saved at K=100,000 in these fixtures.
   
   <details>
   <summary>All 24 scenarios: mean query times and data-byte savings</summary>
   
   P = scan partitions; C = BIGINT payload columns. Negative time change means 
faster. Both time comparisons are shown: versus fused isolates reader pruning; 
versus unfused includes fusion's cost.
   
   | Layout | P | C | K | Unfused ms | Fused ms | Pruning ms | Time vs fused | 
Time vs unfused | Scan bytes saved |
   |---|---:|---:|---:|---:|---:|---:|---:|---:|---:|
   | ascending | 1 | 1 | 16 | 67.5 | 54.5 | 43.5 | -20.2% | -35.6% | 81.6% |
   | ascending | 1 | 1 | 100,000 | 110.5 | 107.5 | 89 | -17.2% | -19.5% | 81.6% 
|
   | ascending | 1 | 16 | 16 | 230.5 | 214.5 | 76 | -64.6% | -67.0% | 77.6% |
   | ascending | 1 | 16 | 100,000 | 346.5 | 355 | 207 | -41.7% | -40.3% | 77.6% 
|
   | ascending | 4 | 1 | 16 | 70.5 | 59.5 | 46 | -22.7% | -34.8% | 82.1% |
   | ascending | 4 | 1 | 100,000 | 320 | 313.5 | 315.5 | +0.6% | -1.4% | 46.3% |
   | ascending | 4 | 16 | 16 | 249 | 231 | 85 | -63.2% | -65.9% | 80.3% |
   | ascending | 4 | 16 | 100,000 | 807.5 | 812 | 686 | -15.5% | -15.0% | 60.7% 
|
   | descending | 1 | 1 | 16 | 168 | 165 | 166.5 | +0.9% | -0.9% | 0.0% |
   | descending | 1 | 1 | 100,000 | 445 | 435.5 | 436.5 | +0.2% | -1.9% | 0.0% |
   | descending | 1 | 16 | 16 | 306.5 | 403 | 401 | -0.5% | +30.8% | 0.0% |
   | descending | 1 | 16 | 100,000 | 642.5 | 765.5 | 761.5 | -0.5% | +18.5% | 
0.0% |
   | descending | 4 | 1 | 16 | 188 | 183.5 | 178.5 | -2.7% | -5.1% | 0.0% |
   | descending | 4 | 1 | 100,000 | 593 | 583 | 575 | -1.4% | -3.0% | 0.0% |
   | descending | 4 | 16 | 16 | 341.5 | 433 | 448 | +3.5% | +31.2% | 0.0% |
   | descending | 4 | 16 | 100,000 | 1,048 | 1,150 | 1,146 | -0.3% | +9.4% | 
0.0% |
   | random | 1 | 1 | 16 | 47.5 | 44 | 41.5 | -5.7% | -12.6% | 8.2% |
   | random | 1 | 1 | 100,000 | 583.5 | 563.5 | 566 | +0.4% | -3.0% | 0.0% |
   | random | 1 | 16 | 16 | 215.5 | 201 | 182 | -9.5% | -15.5% | 10.3% |
   | random | 1 | 16 | 100,000 | 901 | 972.5 | 971 | -0.2% | +7.8% | 0.0% |
   | random | 4 | 1 | 16 | 73.5 | 57.5 | 58 | +0.9% | -21.1% | 5.2% |
   | random | 4 | 1 | 100,000 | 808.5 | 787 | 789.5 | +0.3% | -2.4% | 0.0% |
   | random | 4 | 16 | 16 | 256.5 | 249 | 243 | -2.4% | -5.3% | 1.2% |
   | random | 4 | 16 | 100,000 | 1,344 | 1,417.5 | 1,434 | +1.2% | +6.7% | 0.0% 
|
   
   </details>
   
   <details>
   <summary>All 24 scenarios: exact reader counters</summary>
   
   Both JVM runs produced identical reader counters. Both baselines emitted all 
**1,048,576 rows**, read the same data bytes, and skipped zero row groups. The 
table shows pruning mode's emitted rows and skipped groups.
   
   | Layout | P | C | K | Baseline data bytes | Pruning data bytes | Pruning 
output rows | Row groups skipped / total |
   |---|---:|---:|---:|---:|---:|---:|---:|
   | ascending | 1 | 1 | 16 | 8,465,139 | 1,553,842 | 192,481 | 5 / 6 |
   | ascending | 1 | 1 | 100,000 | 8,465,139 | 1,553,842 | 192,481 | 5 / 6 |
   | ascending | 1 | 16 | 16 | 73,485,975 | 16,482,556 | 235,199 | 4 / 5 |
   | ascending | 1 | 16 | 100,000 | 73,485,975 | 16,482,556 | 235,199 | 4 / 5 |
   | ascending | 4 | 1 | 16 | 8,465,914 | 1,515,816 | 187,748 | 20 / 24 |
   | ascending | 4 | 1 | 100,000 | 8,465,914 | 4,547,408 | 563,244 | 12 / 24 |
   | ascending | 4 | 16 | 16 | 73,493,013 | 14,458,832 | 206,308 | 20 / 24 |
   | ascending | 4 | 16 | 100,000 | 73,493,013 | 28,918,405 | 412,616 | 16 / 24 
|
   | descending | 1 | 1 | 16 | 8,465,135 | 8,465,135 | 1,048,576 | 0 / 6 |
   | descending | 1 | 1 | 100,000 | 8,465,135 | 8,465,135 | 1,048,576 | 0 / 6 |
   | descending | 1 | 16 | 16 | 73,485,972 | 73,485,972 | 1,048,576 | 0 / 5 |
   | descending | 1 | 16 | 100,000 | 73,485,972 | 73,485,972 | 1,048,576 | 0 / 
5 |
   | descending | 4 | 1 | 16 | 8,465,912 | 8,465,912 | 1,048,576 | 0 / 24 |
   | descending | 4 | 1 | 100,000 | 8,465,912 | 8,465,912 | 1,048,576 | 0 / 24 |
   | descending | 4 | 16 | 16 | 73,493,012 | 73,493,012 | 1,048,576 | 0 / 24 |
   | descending | 4 | 16 | 100,000 | 73,493,012 | 73,493,012 | 1,048,576 | 0 / 
24 |
   | random | 1 | 1 | 16 | 8,465,578 | 7,769,839 | 962,405 | 1 / 6 |
   | random | 1 | 1 | 100,000 | 8,465,578 | 8,465,578 | 1,048,576 | 0 / 6 |
   | random | 1 | 16 | 16 | 73,486,406 | 65,932,698 | 940,796 | 1 / 5 |
   | random | 1 | 16 | 100,000 | 73,486,406 | 73,486,406 | 1,048,576 | 0 / 5 |
   | random | 4 | 1 | 16 | 8,466,437 | 8,023,000 | 993,658 | 2 / 24 |
   | random | 4 | 1 | 100,000 | 8,466,437 | 8,466,437 | 1,048,576 | 0 / 24 |
   | random | 4 | 16 | 16 | 73,493,524 | 72,596,667 | 1,035,799 | 3 / 24 |
   | random | 4 | 16 | 100,000 | 73,493,524 | 73,493,524 | 1,048,576 | 0 / 24 |
   
   </details>
   
   Counters come from separate untimed validation executions. `bytes_scanned` 
(requested ranges) equaled `scan_io_data_bytes` (returned data-page bytes) in 
every case; neither measures physical disk/network traffic, and both exclude 
footer/page-index reads in these fixtures. `scan_io_metadata_bytes` was 
**524,288 bytes per file** in every mode (524,288 or 2,097,152 bytes per 
query). There are no Bloom-filter reads in these fixtures.
   
   `row_groups_pruned_statistics`, `page_index_rows_pruned`, and 
`pushdown_rows_pruned` were **zero throughout**. This fixture has one file per 
task and demonstrates within-file pruning through 
`row_groups_pruned_dynamic_filter`; in other workloads, a threshold available 
when opening a later file can instead increment `row_groups_pruned_statistics`.
   
   ### Validation and completed scope
   
   - All **144 mode cases** passed exact result comparisons against Spark, 
native plan/serialization checks, actual partition counts, and attachment 
checks. Every pruning execution attached one filter per partition with zero 
skipped attachments.
   - Fresh native release build and full Maven reactor install passed; the 
packaged JNI library's SHA-256 matched the release library. Both benchmark JVMs 
exited 0 and reported `BUILD SUCCESS`. Maven's `exec:java` runner emitted a 
Hadoop shutdown-hook classloader warning after completion; there were no 
benchmark execution failures.
   - The feature supports **one direct signed integer sort key over eligible 
native Parquet scans**, with fresh threshold state per execution/task. 
Unsupported shapes and unsafe reader cases retain ordinary reading. Both 
configuration flags remain **experimental and opt-in**.
   
   This completes the reader-pruning scope of the revised four-PR plan. Broader 
eligibility, fusion performance improvements, and additional TopK 
instrumentation can be separate follow-ups.
   
   ### Reproduce
   
   Use a full JDK 21 with its `bin` directory on `PATH`, Rust, and the pinned 
revision above:
   
   ```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
   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 pruning,fused,unfused
   ```
   
   These runs used the equivalent Maven benchmark invocation from the Makefile 
after the full release reactor build, with an isolated Maven cache and the 
freshly built JNI library selected explicitly.
   


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