Skip to content

CometNativeScanExec evaluates the inner (unrewritten) dynamic-pruning subquery from outputPartitioning: q5 fails with SubqueryAdaptiveBroadcastExec.execute(), q64 intermittently returns 0 rows at SF1000 #6133

Description

@Neuw84

Describe the bug

TPC-DS q64 returns 0 rows at scale factor 1000 with Comet native execution, where Spark returns 12,185. The same CometNativeScan feeding Spark's own operators (every spark.comet.exec.<operator>.enabled=false, so the scan output goes through ColumnarToRow) returns the correct 12,185 rows, and at SF10 every configuration agrees with Spark. So the scan's values are right when read through the ColumnVector accessors, 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_ui instance of the cross_sales CTE against the broadcast of store_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.

configuration q64 rows
Spark (no Comet) 12,185
Comet scan + Comet native operators + Comet shuffle 0 (at 200 and 300 shuffle partitions)
Comet scan, Spark operators and shuffle (spark.comet.exec.*.enabled=false) 12,185
any of the above at SF10 12,185-equivalent, all agree

Steps to reproduce

  1. Generate TPC-DS SF1000 as parquet (we used the tpcds-kit dsdgen, char/varchar columns stored as strings; SF10 does not reproduce).
  2. Run q64 (the official TPC-DS text, parameters i_color in ('purple','burlywood','indian','spring','floral','medium'), i_current_price between 64 and 74, cs_ui.syear = 1999 / cs2.syear = 2000) with:
--conf spark.plugins=org.apache.spark.CometPlugin
--conf spark.comet.enabled=true --conf spark.comet.scan.enabled=true
--conf spark.comet.exec.enabled=true --conf spark.comet.exec.shuffle.enabled=true --conf spark.comet.exec.shuffle.mode=auto
--conf spark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager
--conf spark.comet.explainFallback.enabled=true --conf spark.comet.cast.allowIncompatible=true
--conf spark.memory.offHeap.enabled=true --conf spark.memory.offHeap.size=16g
--conf spark.sql.shuffle.partitions=300

→ 0 rows.

  1. Same, adding --conf spark.comet.exec.shuffle.enabled=false and spark.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 showing CometNativeScan for 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 formatted output 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.

Activity

  1. Neuw84 commented on Sep 23, 2026

    @Neuw84
    Author

    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 CometNativeScan under 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 reads store_sales through a dynamic-partition-pruning filter on ss_sold_date_sk; in the same environment we also see, on the third-party-operator configuration, q5 fail with SubqueryAdaptiveBroadcastExec does not support the execute() code path — a DPP subquery evaluated before AQE rewrote it. A timing-dependent evaluation of the scan's DPP InSubqueryExec (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.

  2. Neuw84 commented on Sep 23, 2026

    @Neuw84
    Author

    Caller stack captured for the q5 failure (Comet 1.0.0, Spark 4.1.3, CometNativeScan with 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 serializedPartitionData describes: the inner scan.partitionFilters holds a separate InSubqueryExec instance for the dynamic-pruning filter, which Spark's expression walk does not see — so PlanAdaptiveDynamicPruningFilters rewrites the outer copy but the inner one is still the placeholder SubqueryAdaptiveBroadcastExec when outputPartitioning (reached by ValidateRequirements inside optimizeQueryStage) forces serializedPartitionData, which calls updateResult() 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 inner InSubqueryExec resolves 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 ValidateRequirements walks (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 (rewrite scan.partitionFilters when the wrapper's partitionFilters are rewritten, or evaluate only the wrapper's already-resolved InSubqueryExec values) and to keep outputPartitioning from evaluating runtime subqueries at all (originalPlan.outputPartitioning, or a count that does not depend on dynamic pruning).

  3. 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
  4. Neuw84 commented on Sep 24, 2026

    @Neuw84
    Author

    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_datafusion set 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 the SubqueryAdaptiveBroadcastExec placeholder 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.

  5. 5 remaining items

  6. dwsmith1983 commented on Sep 27, 2026

    @dwsmith1983
    Contributor

    Do you need some help with the tests? I can run both to cases.

    Yes, that would help. On the SF1000 setup:

    1. q64 with spark.sql.exchange.reuse=false. If the reuse problem is the cause, it should stop returning 0 rows.
    2. For a run that still returns 0 rows, the final executed plan. A ReusedExchange under the second cross_sales reference that points at the other year's store_sales join 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.

  7. Neuw84 commented on Sep 28, 2026

    @Neuw84
    Author

    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.reuse on 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_sales references keep distinct store_sales scans, 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 separate CometNativeScan store_sales … dynamicpruningexpression(ss_sold_date_sk IN dynamicpruning#1301) → CometSubqueryBroadcast dynamicpruning#1301, [d_date_sk]

    No ReusedExchange collapses the two store_sales ⋈ store_returns branches; the ReusedExchange nodes present are all benign shared dimension broadcasts (date_dim, customer_demographics, household_demographics, income_band, …) and the shared store_returns shuffle. 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 second cross_sales scan's dynamicpruning# id and whether it becomes a ReusedExchange of the first branch.

    Check 3 — Comet built from main + #6268 + #6270 (the one that matters)

    Built in-cluster from apache/datafusion-comet:

    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.

  8. dwsmith1983 commented on Sep 28, 2026

    @dwsmith1983
    Contributor

    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_sales scans each carrying its own pruning filter and no reuse across the two store_sales ⋈ store_returns joins. In a 0-row run on stock, I'd expect a ReusedExchange in place of the second branch's join (or its broadcast) pointing at the first branch's, while the scan-level plan looks the same.

  9. Neuw84 commented on Sep 28, 2026

    @Neuw84
    Author

    I will do some more runs with those patches, will update the comment l!

  10. added
    bugSomething isn't working
    priority:criticalData corruption, silent wrong results, security issues
    area:scanParquet scan / data reading
    and removed on Sep 28, 2026
  11. comphead commented on Sep 28, 2026

    @comphead
    Contributor

    Thanks @Neuw84 and @dwsmith1983 this looks pretty critical.

  12. dwsmith1983 commented on Sep 28, 2026

    @dwsmith1983
    Contributor

    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.

  13. added this to the 1.2.0 milestone on Sep 28, 2026
  14. Neuw84 commented on Sep 30, 2026

    @Neuw84
    Author

    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.

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

Metadata

Metadata

Labels

area:scanParquet scan / data readingbugSomething isn't workingpriority:criticalData corruption, silent wrong results, security issues

Type

No type

Projects

No projects

    Milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions