Repository navigation
fix: do not resolve DPP subqueries when AQE reads a native scan's partitioning - #6268
dwsmith1983 wants to merge 21 commits into
Conversation
…titioning AQE validates a stage by reading outputPartitioning after CoalesceShufflePartitions changes it, which is before Comet converts the DPP placeholders. CometNativeScanExec answered with its post-pruning split count, which ran the SubqueryAdaptiveBroadcastExec placeholder and threw. A non-bucketed scan now reports UnknownPartitioning(0), as FileSourceScanExec does, and CometNativeExec sizes the native RDD from perPartitionData at execution. Part of apache#6133.
…artitioning CometIcebergNativeScanExec counted its tasks in outputPartitioning, which runs BatchScanExec.inputRDD and translateRuntimeFilterV2 on a DPP subquery that has no result yet while AQE validates a stage. It now reports UnknownPartitioning(0), as BatchScanExec does on Spark 3.5 and later, and CometNativeExec sizes the native RDD from the scan's tasks at execution.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: AQE could read scan partitioning before DPP subqueries were executable, causing native Parquet and Iceberg scans to fail during stage validation.
- Design approach: Report
UnknownPartitioning(0)during planning and obtain execution partition counts from resolvedperPartitionData. - Correctness / compatibility analysis: Checked Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0, plus standalone scans, fused execution, broadcast sizing, shuffle inputs and limit consumers. The change matches non-bucketed V1 scans and ordinary V2 scans on Spark 3.5+.
- Key design decisions: Preserve bucketed partitioning and resolve subqueries before reading execution data. Reusing the existing lazy scan data adds no per-row work or unnecessary abstraction.
- Implementation sketch: Update both scan partitioning overrides, add scan-specific sizing cases in
buildNativeContext, and cover coalescing and bucketed joins with three regression tests. - Behavioral changes worth calling out: A single-split native scan no longer advertises an
AllTuplesguarantee through its partition count. Spark 3.4 V2 scans can reportSinglePartition, so the Iceberg override is deliberately more conservative there. - Suggested improvements: No introduced P1/P2 issues found within this review.
Reviewed both commits and all five changed files in the full diff from base 605051ad239ef704f5f25d67910a446a6b6d7c70 to head a3e1ea43e3073c841565c4ed3ec604679c7b4e90. The PR is non-draft. No existing reviews, comments or threads were present. Routed skills: review-comet-pr and review-comet-shuffle-pr for partitioning and execution sizing.
Exact-head CI: the label check passed. Comet CI, CodeQL and the title workflow remain action_required; Comet CI has no executed jobs or build/test verdict.
Validation: make core passed. The Spark 4.1.3 reactor build and 53 focused DPP/plan-data tests passed, including all three new regressions, with no failures or cancellations. git diff --check passed. Other Spark profiles were checked against source only. Full Spark SQL/Iceberg upstream suites and platform matrices were not run; ScalaStyle and Spotless were skipped. The working tree remains clean.
|
@andygrove could you approve a CI run on |
|
We tested this PR together with #6270 on TPC-DS at SF1000: Parquet on S3, a Spark 4.1 cluster on EKS, 8 executors x 13 cores on m5.4xlarge, one availability zone, AQE on with 300 shuffle partitions. The build was main
Thanks for the fix. |
* docs(bench): name the Comet build and the engine comparison from the page meta render-benchmark-page.py hard-coded "DataFusion Comet 1.0.0" in the chart legend, the headline card and the summary table, and a fixed reading of where VecRuntime and Comet each lead. The meta now carries both (comet_label, vector_vs_comet); the defaults are the old text, so a page rendered without the keys is unchanged. Needed for the 2026-09-30 x86 run, whose Comet is a source build with two unmerged fixes. Crew: Butterfly (Kiro Crew Issue Radar) * docs(bench): TPC-DS 1 TB pages from the 2026-09-30 single-AZ runs, x86 with patched Comet Both pages are regenerated from runs made on 2026-09-30, each in one cluster session with all nine nodes in one availability zone (us-east-1b) and S3 through a gateway VPC endpoint, on main cb755d1 (#554), with AQE at its defaults and AOT off. - x86 (9 x m5.4xlarge): Spark 3,313.2 s, VecRuntime 2,419.7 s (1.37x, geomean 1.31x, faster on 88 of 103), Comet 2,177.9 s (1.52x). Comet is a source build of main b58b2f3a with the unmerged apache/datafusion-comet#6268 and #6270: q5 completes and q64 returns Spark's 12,185 rows, so every engine now matches Spark's results except q65 (ties). - Graviton4 (9 x m8g.4xlarge): Spark 2,470.2 s, VecRuntime 1,892.7 s (1.31x, geomean 1.31x, faster on 96 of 103); the x86 reference on the page is the run above. render-benchmark-page.py: the TL;DR's closing sentence is a meta key (results_line); its default is the old text, so a page rendered without it is byte-identical. The README headline table and the Graviton line carry the new numbers; docs/results.md records both runs. Crew: Butterfly (Kiro Crew Issue Radar) --------- Co-authored-by: Butterfly (Kiro Crew) <butterfly-crew@users.noreply.github.com>
|
I would like to do a deeper review of this PR before it is merged. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: AQE requirement validation could force unresolved DPP subqueries through native scan partition counting, causing Parquet and Iceberg queries to fail.
- Design approach: Report
UnknownPartitioning(0)during planning and obtain execution partition counts from resolved scan data. - Correctness / compatibility analysis: Checked Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Traced standalone scans, fused execution, broadcast alignment, shuffle inputs, unions and limit consumers. Execution retains the actual task count.
- Key design decisions: Preserve bucketed Parquet partitioning and resolve subqueries before accessing scan data. Reusing existing lazy data adds no per-row work or unnecessary abstraction.
- Implementation sketch: Update both scan partitioning overrides, add scan-specific sizing in
buildNativeContext, and add three regression tests covering coalescing, pruning and bucketed joins. - Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, the intended changes avoid premature DPP execution and stop advertising anAllTuplesguarantee for single-split non-bucketed scans. Spark 3.4 V2 scans can reportSinglePartition, so the Iceberg override is deliberately more conservative there. - Suggested improvements: No introduced P1/P2 issues found within this review.
Reviewed all five files in the full diff from base fef94f6cd78b18151dff57b7a936798385356de5 to head 6a074c6b3fd62295a6b2dd802a2c26f2959a43bd. The PR is non-draft. Read the supplied reviews, issue comments and empty inline/thread collections. No substantiated existing P1/P2 concern remains unresolved. Routed skills: review-comet-pr, review-comet-shuffle-pr and review-comet-ffi-pr.
Exact-head CI: the label check passed. Comet CI and CodeQL report action_required, providing no build/test verdict.
Validation: a fresh native build and Spark 4.1.3 reactor run passed all 32 TopK and native-shuffle input tests. Verified existing exact-head reactor results for 57 DPP and scan-data tests, including all three new regressions. No failures or cancellations. git diff --check passed and the working tree is clean. Other Spark profiles were checked against source only. Full upstream Spark SQL/Iceberg suites and platform matrices were not run. ScalaStyle and Spotless were skipped.
# Conflicts: # spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala
andygrove
left a comment
There was a problem hiding this comment.
I went through this in more depth as I said I would. Without the change, the new coalescing tests fail on 3.5 and 4.1 with the errors from #6133, and the Iceberg one fails on 3.4 too. With it they pass on 3.4, 3.5 and 4.1, and so do Spark's DynamicPartitionPruningSuite and the bucketed read suites with Comet on. Nothing at execution still takes the partition count from outputPartitioning. EnsureRequirements never sees a Comet scan, so the only planning reader is ValidateRequirements, and this can only make that check stricter.
CI hasn't run at the current head, and the Spark SQL and Iceberg suites have never run on this PR. Since this changes what every native scan reports, I've added run-spark-4.1-tests, run-iceberg-tests and run-all-spark-profiles.
While checking the Spark 3.4 SinglePartition case in the description, I found an older bug nearby that this PR doesn't cause. I filed it as #6667.
| * parent's native execution receives an empty input. (`CometIcebergNativeScanExec` does NOT use | ||
| * this trait; it has a dedicated `findAllPlanData` case.) | ||
| * | ||
| * `perPartitionData.length` is the partition count the native block runs with, and the count in |
There was a problem hiding this comment.
Now that buildNativeContext sizes every CometScanWithPlanData from perPartitionData, could this doc also say that outputPartitioning must not read perPartitionData, or anything else that runs the scan's DPP subqueries? AQE reads it while optimizing a stage, before CometPlanAdaptiveDynamicPruningFilters has converted the placeholders, and that's the q5 crash. The Delta scan in #5365 copied the old pattern and works around it with hasUnevaluableSubqueryFilter. With the rule written here it could just report UnknownPartitioning(0) like CometNativeScanExec does now.
There was a problem hiding this comment.
could this doc also say that
outputPartitioningmust not readperPartitionData, or anything else that runs the scan's DPP subqueries?
Added in 2eb4ff6. The CometScanWithPlanData doc now says perPartitionData.length is what sizes the native RDD at execution, that outputPartitioning must not read perPartitionData or anything that evaluates the scan's DPP subqueries because AQE calls it while optimizing a stage before CometPlanAdaptiveDynamicPruningFilters has converted the placeholders, and that a scan that cannot report a real partitioning without them should report UnknownPartitioning(0), as CometNativeScanExec does for a non-bucketed scan.
I pushed this while the label runs were still going so the doc change would not wait on them. It only touches that comment, so the code those runs are testing is unchanged.
…subqueries AQE reads outputPartitioning while it optimizes a stage, before the adaptive DPP placeholders are converted, so a scan that computes it from perPartitionData runs a placeholder subquery. Say so in the CometScanWithPlanData doc, and point a scan that cannot report a real partitioning otherwise at UnknownPartitioning(0), as CometNativeScanExec does for a non-bucketed scan.
A non-bucketed CometNativeScanExec now reports UnknownPartitioning(0), as FileSourceScanExec does, so the coalesced-scan metrics test cannot read the split count from outputPartitioning. perPartitionData holds the count the native RDD is sized from, which is what the precondition meant to check.
|
The failure on every profile was @andygrove can you add it back to the merge queue? It should be fine now. |
Which issue does this PR close?
Part of #6133. This fixes the q5 failure. The q64 row loss has a separate cause, tracked in #6264.
Rationale for this change
AQE applies Spark's own query stage rules before Comet's. When
CoalesceShufflePartitionschanges a stage,optimizeQueryStagerunsValidateRequirements, which readsoutputPartitioningon every node in the stage.CometPlanAdaptiveDynamicPruningFiltershas not run yet at that point, so a scan's DPP subquery is still Spark'sSubqueryAdaptiveBroadcastExecplaceholder.For a non-bucketed scan,
CometNativeScanExec.outputPartitioningreturnedUnknownPartitioning(perPartitionData.length). Computing that length runsInSubqueryExec.updateResulton the placeholder, which throwsSubqueryAdaptiveBroadcastExec does not support the execute() code path. That is the stack trace in the issue.Coalescing skips shuffle reads that share a subtree with a scan, and
CoalesceShufflePartitionsonly splits a stage into separate groups at Spark'sUnionExec,BroadcastHashJoinExec,BroadcastNestedLoopJoinExecandCartesianProductExec. So the failure needs one of those in the scan's stage. q5 has a union and hits it when Comet's union is not used.FileSourceScanExecdoes not have this problem: it reportsUnknownPartitioning(0)for a non-bucketed scan and only resolves DPP when it builds its RDD.CometIcebergNativeScanExecfails the same way. ItsoutputPartitioningcounted the scan's tasks, which goes throughBatchScanExec.inputRDDandtranslateRuntimeFilterV2and throwsCan't translate ... IN dynamicpruning ... no subquery resultbecause the subquery has no result yet. On Spark 3.5 and later,BatchScanExecwithout a key-grouped partitioning reportsUnknownPartitioning(0).What changes are included in this PR?
CometNativeScanExec.outputPartitioningreportsUnknownPartitioning(0)for a non-bucketed scan, asFileSourceScanExecdoes. A bucketed scan still reports its bucket layout.CometIcebergNativeScanExec.outputPartitioningreportsUnknownPartitioning(0).CometNativeExecsizes the native RDD of a block whose leaf is aCometScanWithPlanDataor aCometIcebergNativeScanExecfromperPartitionData.length, at execution, afterfindAllPlanDatahas resolved the scan's subqueries. It used to read that count fromoutputPartitioning. A bucketed scan has one entry per bucket, so its count is unchanged.Planning now sees 0 partitions for a non-bucketed native scan and for an Iceberg native scan, the same as for Spark's scans. The one exception is Spark 3.4, where
BatchScanExecreportsSinglePartitionfor a single input partition and the Iceberg native scan reportsUnknownPartitioning(0). In particular a single split scan no longer satisfiesAllTupleswhile planning, which also matches Spark. The execution paths that need a count take it from the RDD.How are these changes tested?
Two tests in
CometExecSuite, next to the existing AQE DPP tests:UnionExeckeeps the aggregate's shuffle read in the scan's stage. On main it fails with the stack trace from the issue. With the change it matches Spark, and on Spark 3.5 and later it checks that the DPP scan and a coalesced shuffle read share the final stage, and that DPP still prunes to 1 of 10 partitions.CometIcebergNativeSuitegets the same shuffle coalescing test for the Iceberg native scan. Before the Iceberg change it fails with thetranslateRuntimeFilterV2error on Spark 3.5 and 4.1. With it, the answer matches Spark, the Iceberg scan and a coalesced shuffle read share the final stage, and the scan plans 1 of the 10 file tasks.CometExecSuite,CometScanWithPlanDataSuite,CometTopKSuite,CometJoinSuite,CometIcebergNativeSuiteand the DPP fallback suites pass on Spark 3.5. The AQE DPP tests, the Iceberg DPP tests andCometTopKSuitepass on Spark 4.0 and 4.1, and the AQE DPP and Iceberg DPP tests pass on Spark 3.4.