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>>
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
Observations
EXISTSand a correlatedNOT EXISTSwith an inequality predicate onl_suppkey. Spark rewrites these into semi and anti joins with a join filter.LeftAntisort merge join failing on TPC-H q21) and perf: allocation accounting wrapper costs ~5% on TPC-H Q21 from thread-local lookups in the dlopened library #6165 (closed, allocation accounting overhead of about 5% on TPC-H q21). It is not known whether the build used here includes those fixes.explodeoverarray<struct>columns. Iceberg scan falls back to Spark on IS NULL/IS NOT NULL over list/map columns (stale complex-type check) #5731 (closed) describes Iceberg scans falling back to Spark forIS NULL/IS NOT NULLpredicates on list columns, which Spark pushes down forexplode. It is not known whether that applies here.Expected behavior
Comet should be at least as fast as Spark on this query.
Query
Schema of the tables used
Setup
array<struct<...>>column named<table>_data. The queries are the TPC-H queries rewritten to useexplode(<table>_data)in subqueries.zstdcompression and 512 MB row groups.lineitemis partitioned byl_shipdate,ordersbyo_orderdateandpartbyp_brand.mainyet.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.6andspark.memory.storageFraction=0.2.Not verified yet
main.spark.comet.explainFallback.enabled=truewas set for the runs, so fallback reasons should be available (not yet reviewed).