andygrove opened a new issue, #2499: URL: https://github.com/apache/datafusion-ballista/issues/2499
I ran the full TPC-H suite at SF1000 against the latest `main` (32ceaf8, "fix(scheduler): declare distributed EXPLAIN output columns not-null (#2494)"), which now uses DataFusion 55.1.0. I compared it with an earlier Ballista run from Sep 8 (DataFusion 55.0.0), and with Spark 3.5 + [DataFusion Comet](https://github.com/apache/datafusion-comet) `main` (bc4be39) on the same cluster and data. Summary: - Ballista `main` is about **6% faster** than the Sep 8 build overall (696.1s → 652.5s). Most of that comes from cheaper job planning; execution time is roughly unchanged. q05 and q18 (and to a lesser degree q07) are slower in execution. - Spark + Comet completes the suite in 450.0s, so Ballista is currently about **1.45x slower** overall. Ballista is faster on q08, q09 and q17. - **Job planning is 15% of Ballista's total time** (100.2s of 652.5s), 3–8s per query for anything that scans `lineitem` or `orders`. This is #2497: the scheduler re-reads every Parquet footer on every job. Leaving out planning, Ballista's execution time is 546.3s, **1.21x** Spark + Comet. ### Setup All runs use TPC-H SF1000, Parquet (zstd), read from the same S3 location, on the same Kubernetes cluster (linux/amd64), with 2 iterations per query and the mean reported. **Ballista** - 1 scheduler, 32 executors, 16 concurrent tasks per executor - Executor pods: 8 CPU request / 16 CPU limit, 32 GiB memory - Upstream `tpch` benchmark binary from `benchmarks/`, built from the same commit as the scheduler/executor **Spark 3.5 + Comet** - 32 executors, `spark.executor.cores=16` - Executor pods: 8 CPU request / 8 CPU limit, 64 GiB heap + 32 GiB off-heap (Comet native memory) + 10 GiB overhead - Comet native execution and Comet columnar shuffle enabled; otherwise default Comet settings The layouts aren't identical. Ballista executors can burst to 16 CPUs but have less memory; Spark executors are capped at 8 CPUs but have much more memory. ### Results (seconds, mean of 2 iterations) | Query | Ballista Sep 8 (DF 55.0.0) | Ballista 32ceaf8 (DF 55.1.0) | Ballista change | of which planning | Spark + Comet bc4be39 | Ballista / Comet | |---|---:|---:|---:|---:|---:|---:| | q01 | 16.8 | 14.7 | 0.87x | 6.2 | 11.2 | 1.31x | | q02 | 30.2 | 23.9 | 0.79x | 0.3 | 22.9 | 1.04x | | q03 | 31.0 | 29.0 | 0.94x | 7.1 | 14.2 | 2.05x | | q04 | 20.5 | 15.2 | 0.74x | 5.8 | 7.9 | 1.92x | | **q05** | 63.6 | 75.3 | **1.18x** | 5.3 | 38.1 | 1.98x | | q06 | 15.3 | 9.3 | 0.61x | 5.0 | 0.8 | 11.53x | | q07 | 52.5 | 56.3 | 1.07x | 5.2 | 14.7 | 3.81x | | q08 | 40.7 | 37.1 | 0.91x | 6.6 | 46.1 | 0.81x | | q09 | 51.9 | 42.3 | 0.81x | 6.5 | 51.8 | 0.82x | | q10 | 56.0 | 45.2 | 0.81x | 6.5 | 27.2 | 1.66x | | q11 | 18.2 | 15.6 | 0.85x | 0.2 | 13.9 | 1.12x | | q12 | 19.6 | 14.4 | 0.74x | 6.6 | 5.7 | 2.54x | | q13 | 14.2 | 15.5 | 1.10x | 2.5 | 9.6 | 1.62x | | q14 | 15.7 | 10.9 | 0.70x | 4.7 | 2.6 | 4.15x | | q15 | 19.3 | 14.0 | 0.72x | 4.2 | 12.2 | 1.15x | | q16 | 17.4 | 14.9 | 0.85x | 0.2 | 8.8 | 1.69x | | q17 | 31.2 | 22.4 | 0.72x | 4.2 | 25.1 | 0.89x | | **q18** | 51.4 | 65.3 | **1.27x** | 6.4 | 44.4 | 1.47x | | q19 | 19.0 | 14.5 | 0.77x | 2.9 | 8.4 | 1.72x | | q20 | 33.5 | 24.1 | 0.72x | 3.0 | 7.1 | 3.37x | | **q21** | 65.7 | 81.8 | **1.25x** | 8.1 | 68.7 | 1.19x | | q22 | 12.4 | 10.8 | 0.87x | 2.6 | 8.7 | 1.24x | | **Total** | **696.1** | **652.5** | **0.94x** | **100.2** | **450.0** | **1.45x** | - "Ballista change" is 32ceaf8 vs Sep 8 (< 1 means faster). - "of which planning" is the part of the 32ceaf8 time spent between the job being queued and starting, from the scheduler event log (`started_at - queued_at`). - "Ballista / Comet" is Ballista 32ceaf8 vs Spark + Comet (> 1 means Ballista is slower). ### Finding: job planning re-reads every Parquet footer (#2497) While looking at these numbers I found that planning alone takes 3–8s for every query that scans `lineitem` or `orders`. The three queries that don't (q02, q11, q16) plan in 0.2–0.3s. Details and a local repro are in #2497. In short, the scheduler builds a fresh `RuntimeEnv` for every job, so the file statistics cache starts empty each time, and every Parquet footer is fetched from S3 again. It matters most for the short queries. Planning is over 40% of the time for q01, q06, q12 and q14, and more than half of q06's gap to Spark + Comet (9.3s vs 0.8s) is planning. Across the suite, fixing it would bring Ballista from 1.45x to roughly 1.21x of Spark + Comet. ### Changes since Sep 8 The Sep 8 run's event log only covers q01–q20, with one iteration each, so this breakdown is rough. For those 20 queries: - Planning dropped from 136.3s to 89.4s, which accounts for most of the overall improvement. - Execution was roughly flat (472.7s → 465.5s, 0.98x). The slower queries are slower in execution, not planning: - **q05**: execution 54.9s → 69.7s (71.2s / 68.2s this run), so this looks real. - **q18**: execution 43.2s → 58.6s. Noisy (69.2s / 48.0s), but both iterations are slower than Sep 8. - **q07**: execution 43.7s → 50.8s (52.3s / 49.3s). - **q21**: 65.7s → 81.8s end to end. It's missing from the Sep 8 event log, so I can't split it. I plan to rerun these with more iterations to confirm. I'm opening this to share the numbers and to ask whether anyone has ideas about the q05/q07/q18 execution slowdowns, or about the largest execution gaps vs Spark + Comet (q06, q07, q14, q20). -- 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]
