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]