Repository navigation
CometNativeScanExec evaluates the inner (unrewritten) dynamic-pruning subquery from outputPartitioning: q5 fails with SubqueryAdaptiveBroadcastExec.execute(), q64 intermittently returns 0 rows at SF1000 #6133
Description
Activity
- added 7 commits that reference this issue
on Sep 23, 2026 Update: the result is intermittent, not deterministic. Two further full runs today of the same two configurations that returned 0 rows this morning (Comet 1.0.0 end to end, and
CometNativeScanunder the third-party operators) both returned the correct 12,185 rows for q64 at SF1000, same image, same data, 300 shuffle partitions; a run yesterday of the end-to-end configuration was also correct. So the loss is run-dependent. q64's affected join readsstore_salesthrough a dynamic-partition-pruning filter onss_sold_date_sk; in the same environment we also see, on the third-party-operator configuration, q5 fail withSubqueryAdaptiveBroadcastExec does not support the execute() code path— a DPP subquery evaluated before AQE rewrote it. A timing-dependent evaluation of the scan's DPPInSubqueryExec(before the rewrite → exception; before the broadcast is ready → an empty set → every partition pruned → 0 rows) would explain both; we are capturing the caller stack of the q5 evaluation now and will add it here.Caller stack captured for the q5 failure (Comet 1.0.0, Spark 4.1.3,
CometNativeScanwith third-party operators):SparkUnsupportedOperationException: SubqueryAdaptiveBroadcastExec does not support the execute() code path. at org.apache.spark.sql.execution.SubqueryAdaptiveBroadcastExec.doExecute(SubqueryAdaptiveBroadcastExec.scala:44) at org.apache.spark.sql.execution.SparkPlan.getByteArrayRdd(SparkPlan.scala:378) at org.apache.spark.sql.execution.InSubqueryExec.updateResult(subquery.scala:134) at org.apache.spark.sql.comet.CometNativeScanExec.$anonfun$serializedPartitionData$1(CometNativeScanExec.scala:176) at org.apache.spark.sql.comet.CometNativeScanExec.serializedPartitionData$lzycompute(CometNativeScanExec.scala:171) at org.apache.spark.sql.comet.CometNativeScanExec.perPartitionData(CometNativeScanExec.scala:240) at org.apache.spark.sql.comet.CometNativeScanExec.outputPartitioning$lzycompute(CometNativeScanExec.scala:131) at org.apache.spark.sql.execution.exchange.ValidateRequirements$.validateInternal(ValidateRequirements.scala:48-50) ... (ValidateRequirements recursion, from AdaptiveSparkPlanExec.optimizeQueryStage)So the mechanism is the one the code comment in
serializedPartitionDatadescribes: the innerscan.partitionFiltersholds a separateInSubqueryExecinstance for the dynamic-pruning filter, which Spark's expression walk does not see — soPlanAdaptiveDynamicPruningFiltersrewrites the outer copy but the inner one is still the placeholderSubqueryAdaptiveBroadcastExecwhenoutputPartitioning(reached byValidateRequirementsinsideoptimizeQueryStage) forcesserializedPartitionData, which callsupdateResult()on it →execute()→ this exception. The same eager, unsynchronised evaluation of the inner copy is the natural explanation for the intermittent q64 zero rows: whatever the innerInSubqueryExecresolves to at that moment is what prunes the file partitions, and it is not the filter AQE ends up executing.Why it shows with third-party operators and not in the pure Comet configuration is presumably the plan shape
ValidateRequirementswalks (our operators declare required distributions the walk checks down to the scan); the root is independent of them. A fix on Comet's side would be to make the inner copy follow the outer one (rewritescan.partitionFilterswhen the wrapper'spartitionFiltersare rewritten, or evaluate only the wrapper's already-resolvedInSubqueryExecvalues) and to keepoutputPartitioningfrom evaluating runtime subqueries at all (originalPlan.outputPartitioning, or a count that does not depend on dynamic pruning).- changed the title
[-]TPC-DS q64 returns 0 rows at SF1000 with native execution; correct with Spark operators over the same CometNativeScan and at SF10[/-][+]CometNativeScanExec evaluates the inner (unrewritten) dynamic-pruning subquery from outputPartitioning: q5 fails with SubqueryAdaptiveBroadcastExec.execute(), q64 intermittently returns 0 rows at SF1000[/+]on Sep 23, 2026 One more data point from the same 1 TB TPC-DS setup (Comet 1.0.0, Spark 3.5, AQE on, scan-only configuration with every Comet operator disabled). With
spark.comet.scan.impl=native_datafusionset explicitly instead of left at the default selection, q5 completes and its result matches Spark's (100 rows, equal checksum) in three of three runs tonight — under the default selection it failed on every run over two days with theSubqueryAdaptiveBroadcastExecplaceholder trace above. q64 is unchanged: still 0 rows against Spark's 12,185 (the stale pruning value path), so the explicit setting changes the planner path q5 takes but not the underlying issue. Happy to run a patched build if it helps.- added a commit that references this issue
on Sep 24, 2026 5 remaining items
Do you need some help with the tests? I can run both to cases.
Yes, that would help. On the SF1000 setup:
- q64 with
spark.sql.exchange.reuse=false. If the reuse problem is the cause, it should stop returning 0 rows. - For a run that still returns 0 rows, the final executed plan. A
ReusedExchangeunder the secondcross_salesreference that points at the other year'sstore_salesjoin would confirm it.
The fixes are up: #6268 for the q5 failure (Parquet and Iceberg native scans) and #6270 for the q64 row loss. A build of main with both applied, run at SF1000, is the check that matters. If q64 still comes back empty with AQE off, the plan from that run would help too.
- q64 with
Ran the checks at TPC-DS SF1000 on an 8-executor Spark 4.1.3 / Scala 2.13 cluster (JDK 25), 300 shuffle partitions, TPC-DS Parquet data partitioned by
ss_sold_date_sk. Each repeat is a separate application (fresh JVM / fresh AQE planning), since the loss is run-dependent. Row counts and checksums below are the runner's own (a checksum over the sorted result rounded to 10 significant digits; identical checksum = identical result).Check 1 — stock Comet 1.0.0, end-to-end, q64,
spark.sql.exchange.reuseon vs off (AQE on)exchange.reuse run rows checksum wall (s) false 1 12185 bce69949 157.6 false 2 12185 bce69949 151.6 false 3 12185 bce69949 144.6 true (default) 1 12185 bce69949 66.1 true (default) 2 12185 bce69949 70.5 true (default) 3 12185 bce69949 70.0 All six returned 12,185 rows with the identical checksum, i.e. the 0-row loss did not reproduce in this batch — neither with reuse off (as expected) nor with reuse on. This is consistent with the intermittency already reported: some days the same image/data returns the correct result on every run.
Check 2 — final executed plan of a 0-row run
No 0-row run occurred in Check 1, so there is no empty-result plan to analyse. For reference, in the (correct, 12,185-row) reuse-on final plan the two
cross_salesreferences keep distinctstore_salesscans, each with its own dynamic-pruning filter:- first
cross_sales:CometNativeScan store_sales…PartitionFilters: [isnotnull(ss_sold_date_sk), dynamicpruningexpression(ss_sold_date_sk IN dynamicpruning#1300)]→CometSubqueryBroadcast dynamicpruning#1300, [d_date_sk] - second
cross_sales: a separateCometNativeScan store_sales…dynamicpruningexpression(ss_sold_date_sk IN dynamicpruning#1301)→CometSubqueryBroadcast dynamicpruning#1301, [d_date_sk]
No
ReusedExchangecollapses the twostore_sales ⋈ store_returnsbranches; theReusedExchangenodes present are all benign shared dimension broadcasts (date_dim,customer_demographics,household_demographics,income_band, …) and the sharedstore_returnsshuffle. That is exactly the correct shape — distinct pruning filters, no cross-branch exchange reuse — which matches the correct row count. When a 0-row run is next captured, the plan to look at is the secondcross_salesscan'sdynamicpruning#id and whether it becomes aReusedExchangeof the first branch.Check 3 — Comet built from
main+ #6268 + #6270 (the one that matters)Built in-cluster from
apache/datafusion-comet:- base
main=786bbc8fd630abfd0b29bdaffea6168f9f228863 - fix: do not resolve DPP subqueries when AQE reads a native scan's partitioning #6268 head =
a88941e2ed7ba2bbc02bb687eaed11b8fa1663fa - fix: keep unconverted DPP filters in a scan's canonical form #6270 head =
48708be48e3032a9602009e005aa1238316a2e68 - merge commit =
c2689fb1b9171493f389791f6db4397160536a76(6 commits over base; touchesCometScanUtils.scala,operators.scala) - profile
-Pspark-4.1 -Pscala-2.13, nativelibcomet.sobuilt-Ctarget-cpu=x86-64-v3.
AQE on, exchange reuse on, default settings unless noted. Baselines are plain Spark on the same data.
configuration query AQE runs rows (each) checksum matches Spark spark (baseline) q5 on 1 100 c7eda84b — spark (baseline) q64 on 1 12185 bce69949 — Comet end-to-end q5 on 3 100 c7eda84b yes Comet end-to-end q64 on 3 12185 bce69949 yes Comet end-to-end q64 off 1 12185 bce69949 yes Comet scan under third-party operators q5 on 3 100 c7eda84b yes Comet scan under third-party operators q64 on 3 12185 bce69949 yes Comet scan under third-party operators q64 off 1 12185 bce69949 yes Approximate wall times: Comet end-to-end q5 ≈ 36 s, q64 ≈ 69 s; third-party-operator config q5 ≈ 32 s, q64 ≈ 64 s (q64 AQE-off ≈ 108 s).
Results with the two fixes applied:
- q64 returns 12,185 rows on every run in both configurations, with AQE on and with AQE off — no row loss.
- q5 completes with 100 rows and a checksum equal to Spark's in both configurations — including the "Comet scan under the third-party operators" configuration, which is where q5 previously failed with
SubqueryAdaptiveBroadcastExec does not support the execute() code path. That failure did not recur.
So with #6268 + #6270, q5 and q64 match Spark in both configurations, with AQE on and off. For q5 that is a clear change: under the default scan selection it failed on every run we made over two days. For q64 it is weaker evidence than it looks, because stock 1.0.0 did not lose rows in this batch either (Check 1), so these runs cannot tell the fix apart from a good day. We will keep running q64 on stock with reuse on and post the plan of the next 0-row run.
- first
Thanks for running all three. q5 is the clear result: the failure under the default scan selection is gone in both configurations, and the checksums match Spark.
Agreed on q64: a clean batch on stock can't separate the fix from a good day. Your reuse-on plan has the shape #6270 keeps, with two separate
store_salesscans each carrying its own pruning filter and no reuse across the twostore_sales ⋈ store_returnsjoins. In a 0-row run on stock, I'd expect aReusedExchangein place of the second branch's join (or its broadcast) pointing at the first branch's, while the scan-level plan looks the same.I will do some more runs with those patches, will update the comment l!
- addedbugSomething isn't workingSomething isn't workingpriority:criticalData corruption, silent wrong results, security issuesData corruption, silent wrong results, security issuesarea:scanParquet scan / data readingParquet scan / data readingand removed
on Sep 28, 2026 Thanks @Neuw84 and @dwsmith1983 this looks pretty critical.
this looks pretty critical.
The two fixes are up: #6268 for the q5 failure (Parquet and Iceberg native scans) and #6270 for the q64 row loss, filed as #6264. Both are waiting on CI approval. Neuw84's SF1000 run with both applied is in the comment above: q5 passes in both configurations, and q64 matches Spark on every run.
Running a full test now, will post the results in few hours.
q5 now completes. 23.1 s, 100 rows, the same checksum as vanilla Spark, with 77/77 operators on Comet (11x CometNativeScanExec). Before this fix, q5 failed at this scale with native scans.
The whole run passes. All 103 queries complete. Row counts match vanilla Spark on every query, and checksums match on all but q65, which has ties in its ORDER BY and differs the same way for every engine we compare.Seems that this is fixed with those patches.
Describe the bug
TPC-DS q64 returns 0 rows at scale factor 1000 with Comet native execution, where Spark returns 12,185. The same
CometNativeScanfeeding Spark's own operators (everyspark.comet.exec.<operator>.enabled=false, so the scan output goes throughColumnarToRow) returns the correct 12,185 rows, and at SF10 every configuration agrees with Spark. So the scan's values are right when read through theColumnVectoraccessors, and the rows are lost only when the scan's batches are consumed by native operators.Where the rows go: with operators disabled one at a time, the loss sits in the join of the second
cs_uiinstance of thecross_salesCTE against the broadcast ofstore_sales (ss_sold_date_sk in 2000) ⋈ store_returns; before that join the intermediate row counts agree with Spark. Turning AQE off does not change the result. An unrelated third-party columnar operator set that consumes Comet's scan output through the Arrow C data interface loses the rows in exactly the same place, which suggests the difference is in the exported batches rather than in a particular operator.Environment: Comet 1.0.0, Spark 4.1.3 (
spark-4.1.3-bin-hadoop3), Java 25 (Corretto), Kubernetes (Spark operator), TPC-DS SF1000 parquet on S3 read through S3A. 8 executors × 13 cores, 18 GB heap + 16 GB overhead,spark.memory.offHeap.size=16g.spark.comet.exec.*.enabled=false)Steps to reproduce
tpcds-kitdsdgen,char/varcharcolumns stored as strings; SF10 does not reproduce).i_color in ('purple','burlywood','indian','spring','floral','medium'),i_current_price between 64 and 74,cs_ui.syear = 1999/cs2.syear = 2000) with:→ 0 rows.
--conf spark.comet.exec.shuffle.enabled=falseandspark.comet.exec.{project,filter,aggregate,sort,localLimit,globalLimit,takeOrderedAndProject,hashJoin,sortMergeJoin,broadcastHashJoin,broadcastExchange,expand,union,window,coalesce,collectLimit,explode,sample}.enabled=false, default shuffle manager → 12,185 rows, the plan still showingCometNativeScanfor all 48 scans.Expected behavior
12,185 rows, as Spark and as Comet's scan under Spark's operators.
Additional context
I can attach
explain formattedoutput of both plans and the per-join intermediate row counts if useful. Related to the much older #74 (q64 among the queries once disabled for result differences), but here the scan alone is fine and the difference appears only at scale.