Skip to content

Investigate nested TPC-H q21 slowdown: Comet 37% slower than Spark at SF1000 #6467

Description

@comphead

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

Expected behavior

Comet should be at least as fast as Spark on this query.

Query

-- 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
Schema of the tables used
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>>

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions