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]

Reply via email to