You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
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.
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.
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.
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:
Struct-typed scalar subquery result takes the consuming projection off Comet (widened by Spark 4.2 MergeSubplans) #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.
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.
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.
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.
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.
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.
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: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 tocustomerandcustomer_addresswithca_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 thestore_salesrows 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.
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.
OneRowRelationin Union branches forces Union and downstream aggregates off Comet (TPC-DS q77a) #4949 and Struct-typed scalar subquery result takes the consuming projection off Comet (widened by Spark 4.2 MergeSubplans) #5834 below).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.branch-1.1atee3f239.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.What the plan-stability goldens show
The goldens don't show a fallback that would explain either query:
approved-plans-v1_4/q68, which is fully native (45 of 45 operators, no subqueries or unions).approved-plans-v1_4-spark4_0/q54, which is native except for the SparkSubquerywrappers around its two scalar subqueries ondate_dim.1.0.0andbranch-1.1and 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.1andspark 4.2, and searched for others about operators or expressions that fall back only on 4.x. These could plausibly affect these queries:MergeScalarSubqueriestoMergeSubplansand 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 ondate_dim. The golden shows them unmerged because they group by different expressions, but that should be confirmed in the SF1000 plan.OneRowRelationin Union branches forces Union and downstream aggregates off Comet (TPC-DS q77a) #4949: Spark 4.2 plansOneRowRelationintoUnionbranches, 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.None of these obviously matches q68.
Suggested investigation
spark.sql.ansi.enabled=false.