fix: keep unconverted DPP filters in a scan's canonical form - #6270
dwsmith1983 wants to merge 4 commits into
Conversation
AQE canonicalizes a query stage from its exchange before the stage rules convert DPP placeholders. Comet dropped a DPP filter from the scan's canonical form while it was still a placeholder, so exchanges over scans with different DPP filters compared equal and AQE reused one for both, losing the other branch's rows. The shared helper now drops only DynamicPruningExpression(TrueLiteral), as FileSourceScanExec does. CometNativeScanExec.doCanonicalize drops every DPP filter from originalPlan, a stale copy no DPP rule rewrites, so an unused DPP scan still matches its twin. Closes apache#6264. Part of apache#6133.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Canonicalization discarded unresolved dynamic partition pruning filters, allowing AQE to reuse exchanges across differently pruned scans and lose rows.
- Design approach: Preserve unresolved DPP in the shared scan helper and remove stale DPP copies only from
CometNativeScanExec.originalPlanduring canonicalization. - Correctness / compatibility analysis: The helper matches Spark’s V1 and V2 scan behavior in 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. AQE retains pre-optimization stage identities and separately checks optimized plans for reuse, supporting both parts of this fix.
- Key design decisions: Top-level
partitionFiltersretain the active pruning identity.originalPlanstill contributes schema, projection and other scan identity. The change simplifies the shared helper without adding an abstraction. - Implementation sketch: Two production changes plus three Parquet regression cases and one Iceberg case cover parent broadcast, parent shuffle and same-stage joins.
- Behavioral changes worth calling out: Differently pruned branches retain separate exchange identities. Some parent exchanges remain separate even when pruning later becomes unused, matching Spark’s conservative behavior. No runtime benchmark was performed.
- Suggested improvements: None meeting the P1/P2 reporting threshold. No introduced P1/P2 issues found within this review.
Reviewed the complete four-file diff from 605051ad239ef704f5f25d67910a446a6b6d7c70 to ae530c268d6a7617f851308143a1a917bec6764c. The PR remains non-draft. Snapshot and live discussions contain no reviews, issue comments or review threads. Routed skill: .ai/skills/review-comet-pr/SKILL.md. No specialized sibling skill applies.
Validation: Compiled the exact-head helper and Comet adaptive placeholder class against Spark 4.1.3. A standalone probe confirmed distinct scan and parent-exchange keys for different DPP filters, reproduced their collapse with the base helper, and verified reuse after DPP becomes TrueLiteral. Spark baseline executions of the three added Parquet query shapes returned 400, 28 and 400 rows. These checks do not execute Comet’s native pipeline.
Exact-head CI: Comet CI, CodeQL and the title workflow report action_required, with no executed checks. Only labeling passed. This checkout lacks built Comet native artifacts, so full Comet, Spark SQL and Iceberg suites were not run locally. Their integration verdict remains outstanding.
|
@andygrove could you approve a CI run on |
Which issue does this PR close?
Closes #6264. Part of #6133 (the q64 row loss).
Rationale for this change
CometScanUtils.filterUnusedDynamicPruningExpressionsdropped a DPP filter from a scan's canonical form when its subquery was still the adaptive placeholder, not only when it had becomeTrueLiteral. AQE canonicalizes a query stage from its exchange as it was before the stage optimizer rules ran, which is beforeCometPlanAdaptiveDynamicPruningFiltersconverts the placeholder. So an exchange above that stage saw a scan with no DPP filter. Two such exchanges over scans of the same table with different DPP filters compared equal, and AQE reused the first for both. One branch then read the other's rows. TPC-DS q64 builds the samestore_salesjoin for 1999 and 2000, which is where #6133 loses its rows.The extra stripping came with #4112 so that otherwise identical scans could share a stage while their DPP was still unconverted. That reuse does not need it. The quick
stageCachelookup uses the exchange before optimization, butcreateQueryStageschecks the cache again withnewStage.plan.canonicalizedafter the stage rules have run, and by then the unused filter isTrueLiteraland is dropped as in Spark.What changes are included in this PR?
filterUnusedDynamicPruningExpressionsdrops onlyDynamicPruningExpression(TrueLiteral), matchingFileSourceScanExec. It is shared byCometNativeScanExec,CometScanExecandCometIcebergNativeScanExec.CometNativeScanExec.doCanonicalizedrops every DPP filter fromoriginalPlan. No DPP rule rewritesoriginalPlan, so its filters are a stale copy. The scan's DPP identity is in its top-levelpartitionFilters. Without this, the SPARK-32509 test ("unused DPP filter and exchange reuse") stops reusing, because the stale placeholder inoriginalPlankeeps the scan apart from its twin.One reuse goes away, as in Spark: a parent exchange over two scan stages whose DPP later becomes
TrueLiteralis no longer shared, because its key was fixed while the placeholder was still there. On TPC-DS at SF1 with AQE (andAQEPropagateEmptyRelationexcluded, so queries that return no rows on this data keep their plan shape), the final plans of q14a, q14b, q23a, q23b, q24a and q24b have the same number ofReusedExchangeandReusedSubquerynodes before and after this change. q64 has one moreReusedExchangeand keeps bothstore_salesDPP scans, where main drops one of them.How are these changes tested?
New tests in
CometExecSuiterun aUNION ALLof one CTE joined to two different dimension filters, with AQE on and coalescing off, and check the answer against Spark. On Spark 3.5 and later, where these scans run AQE DPP in Comet, they also check that each branch keeps its own DPP scan pruned to 1 and 3 partitions and that noReusedExchangesits over a DPP scan:CometIcebergNativeSuitegets the q64 shape for the Iceberg native scan. Main returns 100 rows where Spark returns 400 on Spark 3.5 and 4.1, and a wrong answer on 3.4. On 3.4 the Iceberg scan stays in Comet with Spark's DPP rule planning its filter, so this path was exposed there too.With the change, the new tests pass on Spark 3.4, 3.5, 4.0 and 4.1. On 3.5,
CometExecSuite, the DPP fallback suites, the TPC-DS plan stability suites andCometIcebergNativeSuitepass. Reverting the helper change fails the new tests. Reverting theoriginalPlanchange fails the SPARK-32509 test.