andygrove opened a new issue, #6418:
URL: https://github.com/apache/datafusion-comet/issues/6418

   ### Describe the bug
   
   The TPC-DS SF1000 results added in #6308 (Comet 1.1.0 on Spark 4.2.0) show 
Comet slower than Spark on q54 and q68. The previous results were Comet 1.0.0 
on Spark 3.5.8. Times are the mean of two iterations, from 
`benchmarks/results/{1.0.0,1.1.0}/*-tpcds.json`:
   
   | Query | Spark 3.5.8 | Comet 1.0.0 | Speedup | Spark 4.2.0 | Comet 1.1.0 | 
Speedup |
   | --- | ---: | ---: | ---: | ---: | ---: | ---: |
   | q54 | 3.59s | 7.28s | 0.49x | 2.54s | 5.91s | 0.43x |
   | q68 | 2.81s | 1.42s | 1.98x | 2.38s | 3.94s | 0.60x |
   | q46 (control) | 6.62s | 2.82s | 2.35x | 2.82s | 1.77s | 1.59x |
   | q79 (control) | 4.16s | 2.10s | 1.98x | 2.75s | 1.49s | 1.84x |
   
   These look like two separate problems.
   
   **q68 is new.** Comet's time went from 1.42s to 3.94s, while Spark's went 
from 2.81s to 2.38s. No other query got more than 13% slower under Comet 
between the two runs; the next largest are q1 (+13%), q10 (+12%) and q58 
(+11%). q46 has the same shape as q68: the same five-table star join on 
`store_sales`, grouped by ticket, then joined back to `customer` and 
`customer_address` with `ca_city <> bought_city`. Under Comet, q46 got faster 
(2.82s to 1.77s). q68's date filter (`d_dom between 1 and 2`, about 72 days) 
keeps about a quarter of the `store_sales` rows that q46's filter (`d_dow in 
(6,0)`) keeps, so q68 should be the cheaper of the two. It was cheaper in 1.0.0.
   
   **q54 is long-standing.** Comet was already about 2x slower than Spark in 
the 1.0.0 run. Comet 1.1.0 is faster than 1.0.0 on q54 (7.28s to 5.91s), but 
Spark got faster by more between 3.5.8 and 4.2.0.
   
   ### Steps to reproduce
   
   Run TPC-DS q54 and q68 at SF1000 with Comet 1.1.0 on Spark 4.2.0, using the 
configuration on the [TPC-DS benchmark 
page](https://github.com/apache/datafusion-comet/blob/main/docs/source/contributor-guide/benchmark-results/tpc-ds.md).
   
   ### Expected behavior
   
   Comet should be at least as fast as Spark on both queries, and q68 should be 
back near its 1.0.0 result of about 2x faster than Spark.
   
   ### Additional context
   
   #### What changed between the two runs
   
   Several things changed at once, so the q68 slowdown can't be pinned on Comet 
1.1.0 yet. It may be specific to Spark 4.2.
   
   - Spark 3.5.8 became 4.2.0. Comet's Spark 4.2 support is still experimental, 
and 4.2 plans some queries differently (see #4949 and #5834 below).
   - ANSI mode: neither run sets `spark.sql.ansi.enabled`, so ANSI was off in 
the 3.5.8 run and on in the 4.2.0 run, which is the Spark 4 default.
   - Comet 1.0.0 became `branch-1.1` at `ee3f239`.
   - The 1.1.0 Comet run also turned on fallback logging 
(`spark.comet.logFallbackReasons.enabled`, 
`spark.comet.explain.format=verbose`). The cluster, the S3 data, and the 
executor and memory settings were the same in both runs. Native columnar-to-row 
was off in both: 1.0.0 set it off explicitly, and off is the default in 1.1.0.
   - Each query ran only twice, and the committed JSON keeps only the mean, so 
one slow iteration could explain the q68 number.
   
   #### What the plan-stability goldens show
   
   The goldens don't show a fallback that would explain either query:
   
   - On Spark 4.2, q68 resolves to `approved-plans-v1_4/q68`, which is fully 
native (45 of 45 operators, no subqueries or unions).
   - q54 resolves to `approved-plans-v1_4-spark4_0/q54`, which is native except 
for the Spark `Subquery` wrappers around its two scalar subqueries on 
`date_dim`.
   - The Spark 4.x goldens are generated with ANSI on, so ANSI mode alone 
doesn't take either query off Comet.
   - Phase 0 of #6399 compared the golden plans between `1.0.0` and 
`branch-1.1` and found no lost native coverage.
   
   The goldens come from empty tables with AQE disabled, though, so they don't 
show the SF1000 join strategies, AQE's runtime changes, or runtime bloom 
filters, which need a large scan on the application side.
   
   #### Open issues about Spark 4.x gaps
   
   I went through the open issues labelled `spark 4.0`, `spark 4.1` and `spark 
4.2`, and searched for others about operators or expressions that fall back 
only on 4.x. These could plausibly affect these queries:
   
   - #5834: Spark 4.2 renames `MergeScalarSubqueries` to `MergeSubplans` and 
widens it. The merged subquery returns a struct, Comet doesn't support 
struct-typed scalar subqueries, and the projection that consumes it falls back. 
q54 has two scalar subqueries on `date_dim`. The golden shows them unmerged 
because they group by different expressions, but that should be confirmed in 
the SF1000 plan.
   - #4949: Spark 4.2 plans `OneRowRelation` into `Union` branches, which takes 
q77a's unions and aggregates off Comet. q77a isn't in the benchmark set, but 
it's the same kind of plan change that only appears on 4.2.
   - #4968: the BloomFilter tests are skipped on Spark 4.2. This matters if the 
SF1000 plans use runtime bloom filters.
   - #4967 and #5078: ANSI arithmetic differences on Spark 4.2, and the ANSI 
audit follow-ups. ANSI is now on in the benchmark.
   - #2190: string collation support, a Spark 4.0 feature.
   
   None of these obviously matches q68.
   
   #### Suggested investigation
   
   1. If the driver logs from the 1.1.0 run still exist, check the fallback 
reasons logged for q54 and q68, and get the final AQE plans from the event logs.
   2. Re-run q68, with q46 as a control, at least five times on the same setup 
and record every iteration, to confirm q68 is consistently slow.
   3. Separate the Spark version from the Comet version: run q68 with Comet 
1.1.0 on Spark 3.5.8 and on 4.1, and on Spark 4.2 with 
`spark.sql.ansi.enabled=false`.
   4. Compare Comet's final q68 plan on Spark 4.2 against the fastest 
configuration from step 3. Look at join strategies, AQE changes, DPP and 
runtime filters, and any transitions back to Spark.
   5. For q54, compare Spark's and Comet's per-operator SQL metrics on the same 
Spark version to find the stage where Comet loses time.
   


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