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]
