comphead opened a new issue, #6467:
URL: https://github.com/apache/datafusion-comet/issues/6467
## Summary
On a nested-schema benchmark derived from TPC-H (SF1000, Iceberg), q21 is
the one query where Comet is clearly slower than Spark. The mean over 3
iterations is 518.4 s for Comet versus 377.6 s for Spark (+37.3%, 0.73x), and
the iteration ranges do not overlap. Comet is faster than Spark on most of the
other queries in this benchmark, so this looks specific to q21. q21 is also the
longest query, at 45% of Comet's total runtime across the 22 queries (24% for
Spark).
## Results
| Engine | Iteration 1 (s) | Iteration 2 (s) | Iteration 3 (s) | Mean (s) |
|---|---:|---:|---:|---:|
| Spark | 371.732 | 382.426 | 378.716 | 377.625 |
| Comet | 478.485 | 520.237 | 556.420 | 518.381 |
## Observations
- Comet's iteration times rise from one iteration to the next (478.5, 520.2,
556.4 s) while Spark's stay flat. The cause is not known. Memory pressure,
spilling or native memory growth across iterations are possibilities, but this
is only a hypothesis.
- The query has a correlated `EXISTS` and a correlated `NOT EXISTS` with an
inequality predicate on `l_suppkey`. Spark rewrites these into semi and anti
joins with a join filter.
- Earlier upstream issues mention this query: #861 (closed, filtered
`LeftAnti` sort merge join failing on TPC-H q21) and #6165 (closed, allocation
accounting overhead of about 5% on TPC-H q21). It is not known whether the
build used here includes those fixes.
- The query uses `explode` over `array<struct>` columns. #5731 (closed)
describes Iceberg scans falling back to Spark for `IS NULL` / `IS NOT NULL`
predicates on list columns, which Spark pushes down for `explode`. It is not
known whether that applies here.
## Expected behavior
Comet should be at least as fast as Spark on this query.
## Query
```sql
-- using default substitutions
select
supplier_data.s_name,
count(*) as numwait
from
(select s_suppkey, s_nationkey, explode(supplier_data) as supplier_data
from supplier) as supplier,
(select l_orderkey, l_partkey, l_suppkey, l_shipdate,
explode(lineitem_data) as lineitem_data from lineitem) as l1,
(select o_orderkey, o_custkey, o_orderdate, explode(orders_data) as
orders_data from orders) as orders,
(select n_nationkey, n_regionkey, explode(nation_data) as nation_data
from nation) as nation
where
s_suppkey = l1.l_suppkey
and o_orderkey = l1.l_orderkey
and orders_data.o_orderstatus = 'F'
and l1.lineitem_data.l_receiptdate > l1.lineitem_data.l_commitdate
and exists (
select
*
from
lineitem l2
where
l2.l_orderkey = l1.l_orderkey
and l2.l_suppkey <> l1.l_suppkey
)
and not exists (
select
*
from
(select l_orderkey, l_partkey, l_suppkey, l_shipdate,
explode(lineitem_data) as lineitem_data from lineitem) as l3
where
l3.l_orderkey = l1.l_orderkey
and l3.l_suppkey <> l1.l_suppkey
and l3.lineitem_data.l_receiptdate >
l3.lineitem_data.l_commitdate
)
and s_nationkey = n_nationkey
and nation_data.n_name = 'SAUDI ARABIA'
group by
supplier_data.s_name
order by
numwait desc,
supplier_data.s_name
limit 100
```
<details>
<summary>Schema of the tables used</summary>
```
supplier
s_suppkey BIGINT, s_nationkey BIGINT,
supplier_data ARRAY<STRUCT<s_name STRING, s_address STRING, s_phone
STRING, s_acctbal DECIMAL(12,2), s_comment STRING>>
lineitem (partitioned by l_shipdate)
l_orderkey BIGINT, l_partkey BIGINT, l_suppkey BIGINT, l_shipdate DATE,
lineitem_data ARRAY<STRUCT<l_linenumber INT, l_quantity DECIMAL(12,2),
l_extendedprice DECIMAL(12,2),
l_discount DECIMAL(12,2), l_tax DECIMAL(12,2), l_returnflag STRING,
l_linestatus STRING,
l_commitdate DATE, l_receiptdate DATE, l_shipinstruct STRING, l_shipmode
STRING, l_comment STRING>>
orders (partitioned by o_orderdate)
o_orderkey BIGINT, o_custkey BIGINT, o_orderdate DATE,
orders_data ARRAY<STRUCT<o_orderstatus STRING, o_totalprice DECIMAL(12,2),
o_orderpriority STRING,
o_clerk STRING, o_shippriority INT, o_comment STRING>>
nation
n_nationkey BIGINT, n_regionkey BIGINT,
nation_data ARRAY<STRUCT<n_name STRING, n_comment STRING>>
```
</details>
## Setup
- Benchmark: a nested-schema variant derived from TPC-H at scale factor
1000. Each table keeps its key columns (and partition columns) at the top level
and stores every other column in one `array<struct<...>>` column named
`<table>_data`. The queries are the TPC-H queries rewritten to use
`explode(<table>_data)` in subqueries.
- Tables: Iceberg 1.5.0 (downstream build) with Parquet data files, `zstd`
compression and 512 MB row groups. `lineitem` is partitioned by `l_shipdate`,
`orders` by `o_orderdate` and `part` by `p_brand`.
- Cluster: Spark 3.4.3 (downstream build, Scala 2.13) on Kubernetes with
16-core amd64 executors. The Comet build is a downstream build and has not been
checked against upstream `main` yet.
- Recorded configs (both runs):
`spark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager`,
`spark.comet.exec.shuffle.enabled=true`,
`spark.comet.exec.shuffle.compression.codec=lz4`,
`spark.comet.expression.allowIncompatible=true`,
`spark.comet.explainFallback.enabled=true`, `spark.memory.fraction=0.6` and
`spark.memory.storageFraction=0.2`.
- Method: a Spark baseline run (Comet not active) and a Comet run. Each run
is one Spark application that ran all 22 queries in order, with 3 consecutive
iterations per query. Times are per-iteration query execution times in seconds
as recorded by the benchmark harness. There is a single run per engine.
## Not verified yet
- Not reproduced on upstream `main`.
- No physical plans, Spark UI metrics or fallback reasons are attached yet.
`spark.comet.explainFallback.enabled=true` was set for the runs, so fallback
reasons should be available (not yet reviewed).
- Executor count, executor memory and off-heap sizing are not recorded by
the benchmark harness and are not listed here.
- The harness does not record how Comet was switched off in the Spark
baseline run. Both runs list the Comet shuffle manager in their recorded
configs.
--
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]