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]

Reply via email to