Repository navigation
Conversation
andygrove
left a comment
There was a problem hiding this comment.
Thanks for digging into this. The operator analysis is careful, and the negative controls in the fixtures (reversed build order, WHERE bad = 'A', rn <= 2) are what convinced me the model is right rather than just conservative.
I checked each operator the rule treats as an early-stop point against Spark master. BaseLimitExec codegen guards consume behind the count check and the producer loop tests stopEarly(), so rows past the limit really are never projected. TakeOrderedAndProjectExec uses take(limit) only when SortOrder.orderingSatisfies(child.outputOrdering, sortOrder), and the guard here mirrors Spark's own condition exactly. HashJoin.semiJoin / antiJoin and BroadcastNestedLoopJoinExec.leftExistenceJoin all short-circuit through .exists. ExistenceJoin short-circuits too, but Comet rejects it in all three join serdes, so its condition is already evaluated per row in Spark and there is no gap there. The materializesInput list also looks sound: anything missing from it produces more fallback rather than less, and everything in it does drain its input.
The one I think is wrong is WindowGroupLimit, and I have left a comment inline.
My main concern is scope, and it is the thing I would like settled before this lands. Nothing about the bug is specific to unbase64. Comet evaluates eagerly over a batch while Spark evaluates lazily per row, so every expression that can throw on data has this exact problem under these exact operators. arithmetic_ansi.sql shows Comet raising ARITHMETIC_OVERFLOW natively, so SELECT a + b FROM ansi_int_overflow LIMIT 1 with ANSI on fails in Comet and succeeds in Spark for the same reason, and so do ANSI cast, decimal divide-by-zero, and element_at. The machinery here is general; only the trigger is hard-coded. I have written that up inline along with the flip side, because I think it decides the shape of the PR.
The rest of my comments are independent of how that question lands.
One documentation note that I could not attach inline because the file is not in the diff. docs/source/user-guide/latest/expressions.md:618 still has an empty Notes cell for unbase64, and the audit entry at docs/source/contributor-guide/expression-audits/string_funcs.md:248 still lists failOnError and non-trivial children as the only restrictions. Both are committed files, so GenerateDocs will not pick this up. Could you add the new operator-level restriction to both? The getUnsupportedReasons() string reads well, so it is mostly a matter of mirroring it.
Share fixtures and assertions across the LIMIT, join, window and AQE regressions. Keep unique error, ordering, buffer and native-admission controls while removing duplicate planner permutations and SQL fixtures. Leave production behavior unchanged and reduce PR apache#5533 to about 800 changed lines.
andygrove
left a comment
There was a problem hiding this comment.
The window group limit narrowing, the two doc files, and the executed top-K and shuffled-hash/sort-merge coverage all look resolved to me, and RequiresSparkEvaluationMask is the shape I was hoping for on the trigger.
The performance point from my earlier comment is still open. There is still no way to turn the fallback off, and the description says the cost has not been benchmarked. The reach also grew since I raised it. hasLimitAncestor at CometExecRule.scala:810 is passed down unchanged, so it crosses exchanges and query stages, and with AQE on a LIMIT anywhere above a final aggregate whose result expressions decode and whose buffers are not mixed-execution safe now pushes both halves of that aggregate back to Spark. On top of that, every df.show() and every BI-tool LIMIT over a projection containing unbase64 gives back what the native kernel in #5451 bought, on data that is almost always well formed. Could this be gated, either on a new CometConf key or on the existing COMET_EXPR_ALLOW_INCOMPATIBLE, so a user who knows their base64 is clean can keep the native path?
I also still would like a tracking issue for the rest of the family. CometEvaluationMaskSuite now asserts that ANSI Add and Cast return None from evaluationMaskName, which pins the divergence I described rather than recording it, and WHERE k = 1 AND unbase64(bad) = X'616263' and CASE WHEN k = 1 THEN unbase64(bad) END are still out of reach for this mechanism with no LIMIT involved. Could you file that and link it here?
Two other things worth settling. spark/src/test/resources/sql-tests/expressions/misc/unbase64_operator_masks.sql was dropped in the consolidation. The Scala suite runs the same queries, so coverage is not lost, but that fixture was the readable statement of #5532 and it lived where the project asks expression tests to live. Would you keep a minimal version of it? Separately, main has grown preserveSparkAggregateBuffers at CometExecRule.scala:215 from #5537, which un-converts a native Partial back to Spark for the same buffer-compatibility reason as restoreNativeAggregateBuffers. On rebase those two should probably become one helper rather than two near-identical walks in the same file.
andygrove
left a comment
There was a problem hiding this comment.
No change here since my last comment, so the three things I raised are still open: the fallback has no config gate, there is no tracking issue for the rest of the throwing-expression family, and unbase64_operator_masks.sql is still gone.
One thing has got worse in the meantime. The consolidation point is now three-way rather than two-way. Main has preserveSparkAggregateBuffers from #5537, and #5421 adds revertUnsafePartialAggregates with a shared restoreSparkPartial helper, so on rebase this PR's restoreNativeAggregateBuffers becomes a third near-identical walk over the same aggregate and exchange chain in CometExecRule.scala. Worth folding into the shared helper rather than resolving the conflict three ways. Both CometExecRule.scala and CometExpressionSerde.scala conflict with main now, and the branch is 190 commits behind.
The red check is not yours. The Spark 3.5, JDK 17, Scala 2.13 [expressions] job failed on the Download native library step about 60 seconds in, back on 31 August. A re-run will clear it.
Share fixtures and assertions across the LIMIT, join, window and AQE regressions. Keep unique error, ordering, buffer and native-admission controls while removing duplicate planner permutations and SQL fixtures. Leave production behavior unchanged and reduce PR apache#5533 to about 800 changed lines.
ef2986c to
7df1251
Compare
|
Rebased onto current
Validated on Spark 4.1.3, Scala 2.13.17, and JDK 17. All 110 executed targeted tests passed: Full JVM test compilation, Spotless, Scalastyle, and Other Spark/Scala versions, macOS, and the full repository suite were not run locally. Hosted PR CI is separate from these results. |
andygrove
left a comment
There was a problem hiding this comment.
Thanks for adding the config, filing #6006 and restoring the fixture. Those were the three things left from my last round, and they all look right to me.
The thing that has changed since is #5421, which landed on the 20th. This branch now conflicts with main in CometExecRule.scala, and #5421 renamed allAggsSupportMixedExecution to allAggsSupportNativePartialToSparkFinal, which three of the new call sites still use. It also added revertUnsafePartialAggregates, which runs after transform and restores the native Partial chain under any Spark Final whose buffers are unsafe. Once preserveEvaluationMasks has tagged a Final, that pass looks like it does what the throughExchanges = true path of restoreNativeAggregateBuffers does before conversion. When you rebase, could this PR keep only the early tagging of the Final and leave the restoring to revertUnsafePartialAggregates and restoreSparkPartial, so there's one walk over that chain rather than two? The two AQE aggregate tests in CometEvaluationMaskSuite should show quickly whether anything is lost. I'm happy to approve once that's in.
Share fixtures and assertions across the LIMIT, join, window and AQE regressions. Keep unique error, ordering, buffer and native-admission controls while removing duplicate planner permutations and SQL fixtures. Leave production behavior unchanged and reduce PR apache#5533 to about 800 changed lines.
7df1251 to
c7c70fa
Compare
|
@andygrove Addressed in c7c70fa, rebased onto Removed Both AQE aggregate regressions passed, including the fresh/reused native-plan matrix and the actual join change. The focused evaluation-mask, rule, Celeborn, and SQL tests passed: 106 executed tests on Spark 4.1.3 (107 successful registrations including one Spark-3.5-gated body). Updated the description and requested Spark SQL 4.1 CI coverage with |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Batch execution could decode malformed
unbase64values in rows Spark skips, causing otherwise successful queries to fail. - Design approach:
RequiresSparkEvaluationMasksupplies an opt-in expression policy. The planner preserves Spark row execution around affected early-stop operators. - Correctness / compatibility analysis: Checked decoder, limit, join and window semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Reviewed AQE reuse, aggregate-buffer compatibility and consumed-malformed controls. No introduced P1/P2 issues found within this review.
- Key design decisions: Early aggregate tagging protects buffers before materialization. Reusing
revertUnsafePartialAggregatesandrestoreSparkPartialavoids another restoration walk. Exchanges provide boundaries where native execution can resume. - Implementation sketch: The rule tags affected operators, removes intervening batching transitions and restores native ancestors with logical links. Tests cover execution, repeated planning, AQE, dispatcher modes and native controls.
- Behavioral changes worth calling out: Valid input also receives fallback by default.
spark.comet.exec.preserveEvaluationMasks.enabled=falseprovides the documented opt-out. The early bailout limits extra planning work, but performance cost remains unbenchmarked. Broader throwing-expression and conditional-mask work is tracked in#6006. - Suggested improvements: None meeting the P1/P2 reporting bar. Previously raised concerns are addressed in this head.
Reviewed the full seven-commit, twelve-file diff from dd68a531c413c287c1f5fbe12ade4edb488900e5 to c7c70faf46d7fed8131088e2f999166c7306ffdf. Confirmed the PR is not a draft. Routed skills: review-comet-pr, review-comet-expression-pr and review-comet-shuffle-pr. Read existing reviews, issue comments, inline comments and threads, excluding Copilot.
Exact-head CI: 32 checks passed and 14 were skipped, with no failed or pending checks. CI run passed Required Checks, Linux Comet suites and all seven Spark 4.1 SQL shards. Logs confirm all nine CometEvaluationMaskSuite tests and both restored SQL-fixture dispatcher variants passed. Other Spark runtime profiles and macOS were not covered by this run.
Local validation: git diff --check and dev/ci/check-ci-config.py passed. Maven bootstrap failed with DNS resolution failure for repo.maven.apache.org, preventing local Scala execution. The initial native build lacked JNI headers. A retry using JDK 17 reached the 90-second bound without completing. Runtime validation therefore relies on exact-head CI. Project files remain unchanged.
andygrove
left a comment
There was a problem hiding this comment.
unbase64 keeps aggregate buffers compatible through an actual AQE join change is flaky for me, and I think it needs fixing before this goes into the merge queue. I ran CometEvaluationMaskSuite seven times on this branch merged with current main, across Spark 4.1, 3.5 and 3.4, and that test failed in five of them. Every failure is the BuildRight assertion in checkSuccess with join=INNER, which reports List(BuildLeft) did not equal List(BuildRight).
I don't think Comet is doing anything wrong there. For an inner join either side can be broadcast, so Spark's getBuildSide takes the smaller one. QueryStageExec.computeStats returns None until a stage has materialized, and LogicalQueryStage then falls back to the static estimate. So if the right side's shuffle stage finishes first, the aggregate side is compared on its estimate and can win. With spark.sql.adaptive.logLevel=warn you can see the Comet-disabled reference runs in the same test choose Inner, BuildLeft as well. CI passed once at this head, but the suite runs in the [expressions] bucket, so the merge queue will hit it sooner or later. I've left a suggestion inline.
The consolidation looks right to me, and the tests pin it. Taking revertUnsafePartialAggregates out of the pipeline fails the reused-native-subtree leg of unbase64 chooses compatible aggregate buffers before AQE materialization. Dropping the early Final tag fails both aggregate tests, and the real AQE join run then crashes with an NPE in Collect.deserialize under Spark's ObjectHashAggregateExec. There is one gap left in the same test, also inline.
The branch also conflicts with main in expressions.md again, because #6168 rewrote the start of the unbase64 row. Keeping main's first sentence and appending the evaluation-mask sentences from this PR should resolve it. Once the build-side assertion is deterministic and the compatible-buffer check is in, I'm happy to approve.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Batch execution could decode malformed
unbase64values in rows Spark skips, failing otherwise successful queries. - Design approach:
RequiresSparkEvaluationMaskenrolls expressions in operator-level protection. Affected paths retain Spark row execution around early-stop operators. - Correctness / compatibility analysis: Compared relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Reviewed limits, join conditions, window limits, transitions and AQE reuse. No additional introduced P1/P2 issues found within this review. The previously reported AQE test flake remains unresolved.
- Key design decisions: Early aggregate tagging preserves compatible buffers before AQE materialization. Reusing
revertUnsafePartialAggregatesandrestoreSparkPartialavoids another restoration walk. - Implementation sketch: Cached serde enrollment drives sticky planner tags. The rule removes intervening batching transitions and rebuilds affected native ancestors while retaining logical links.
- Behavioral changes worth calling out: Valid input also receives fallback by default.
spark.comet.exec.preserveEvaluationMasks.enabled=falseprovides the documented opt-out. Performance cost remains unbenchmarked. Broader throwing-expression and conditional-mask gaps are tracked in#6006. - Suggested improvements: Resolve the existing flaky build-side assertion before merging. Require
BuildRightonly forLEFT SEMI, and letINNERassert that AQE selected one broadcast join.
The existing P2 blocker is at spark/src/test/scala/org/apache/comet/serde/CometEvaluationMaskSuite.scala:399 (discussion). With AQE enabled and join=INNER, Spark legitimately chooses BuildLeft, but the test requires Seq(BuildRight) and fails despite correct results. Spark's join-selection and query-stage statistics code supports this concern. The reviewer reported failures in five of seven runs across Spark 3.4, 3.5 and 4.1 on the branch merged with then-current main. The assertion remains unchanged at this head and can fail the CI expressions bucket. This is retained as an existing blocker, not duplicated as a new finding.
Reviewed the full seven-commit, twelve-file diff from dd68a531c413c287c1f5fbe12ade4edb488900e5 to c7c70faf46d7fed8131088e2f999166c7306ffdf. Confirmed non-draft status. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-shuffle-pr. Read existing reviews, issue comments, inline comments and threads, excluding Copilot.
Exact-head CI: 32 checks passed and 14 were skipped, with none failed or pending. Run 35928478979 passed Linux Comet suites and all seven Spark 4.1 SQL shards. Logs confirm all nine evaluation-mask tests and both SQL-fixture dispatcher variants passed in that run. Other Spark runtime profiles and macOS were not covered.
Validation limits: Local git diff --check and dev/ci/check-ci-config.py passed. Maven bootstrap failed with UnknownHostException: repo.maven.apache.org, preventing local runtime validation and repetition of the reported flake. Runtime evidence relies on exact-head CI and the attributed existing reproduction. Project files remain unchanged.
|
This is a light fully automated review since there are so many PRs open. When a filter that decodes sits under a limit, this also takes the Parquet scan out of native execution. Spark's |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Batch execution could decode malformed
unbase64values in rows Spark skips, failing otherwise successful queries. - Design approach:
RequiresSparkEvaluationMaskenrolls expressions in planner protection that preserves Spark row evaluation around early-stop operators. - Correctness / compatibility analysis: Compared relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Reviewed limits, join conditions, window limits, transitions and AQE reuse. No additional introduced P1/P2 issues found within this review. Existing concerns below remain unresolved.
- Key design decisions: Early aggregate tagging protects intermediate buffers before materialization. Reusing
revertUnsafePartialAggregatesandrestoreSparkPartialavoids another restoration walk. - Implementation sketch: Cached serde enrollment drives persistent planner tags. Affected native ancestors are rebuilt with logical links, and intervening batching transitions are removed.
- Behavioral changes worth calling out: Valid input also receives fallback by default.
spark.comet.exec.preserveEvaluationMasks.enabled=falseprovides the documented opt-out. Performance cost remains unbenchmarked. Broader throwing-expression and conditional-mask work is tracked in #6006. - Suggested improvements: Address the existing AQE test flake and unnecessary Parquet scan fallback described below.
Existing P2 concerns, retained without duplicate findings:
spark/src/test/scala/org/apache/comet/serde/CometEvaluationMaskSuite.scala:399still requiresSeq(BuildRight)for an inner join. Spark legitimately choosesBuildLeftbased on available stage statistics. The reviewer reported five failures in seven runs across Spark 3.4, 3.5 and 4.1 on the branch merged with then-current main. Spark's join-selection and query-stage statistics sources support that explanation. RequireBuildRightonly forLEFT SEMIand assert one broadcast join forINNER. Existing discussion.spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:968treats scandataFiltersas decoding sites. ForSELECT bad FROM t WHERE unbase64(bad) <=> X'616263' LIMIT 1, this restores the Parquet scan to Spark even with the defaultspark.comet.parquet.rowFilterPushdown.enabled=false. Source tracing confirms Spark retains that predicate indataFilters, while native per-row evaluation is disabled and pruning rejects the column function. The Spark filter already preserves row evaluation. Keep the scan native when it does not evaluate the decoder. The acceleration loss is deterministic, though its runtime magnitude was not measured. Existing discussion.
Reviewed the full seven-commit, twelve-file diff from dd68a531c413c287c1f5fbe12ade4edb488900e5 to c7c70faf46d7fed8131088e2f999166c7306ffdf. Confirmed non-draft status. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-shuffle-pr. Read existing reviews, issue comments, inline comments and threads, excluding Copilot.
Exact-head CI: 32 checks passed and 14 were skipped, with none failed or unfinished. Run 35928478979 passed Linux Comet suites and all seven Spark 4.1 SQL shards. Logs confirm all nine evaluation-mask tests and both SQL-fixture dispatcher variants passed. Other Spark runtime profiles and macOS were not covered.
Validation limits: Local git diff --check and dev/ci/check-ci-config.py passed. Maven bootstrap failed with UnknownHostException: repo.maven.apache.org, preventing local runtime validation and repetition of the reported flake. Runtime evidence relies on exact-head CI and the attributed existing reproduction. Project code remains unchanged.
c7c70fa to
1ad7854
Compare
|
Rebased and addressed the remaining requests in The scan guard now ignores V1 data filters when native row-filter pushdown is disabled, retaining the native scan. The same fixture asserts one native scan; the enabled row-filter control asserts fallback. The AQE test requires one broadcast join for INNER and BuildRight only for LEFT SEMI, and pins All 9 focused suite tests pass on Spark 4.1.3, with current-source JVM compilation and the verified exact-base native library. Earlier policy/config/docs/window/top-K/join requests remain implemented. Fresh hosted CI is separate. |
|
CI follow-up |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Batch execution could evaluate malformed
unbase64values on rows Spark skips, failing otherwise successful queries. - Design approach: Serde enrollment identifies protected expressions. The planner retains Spark row execution around limits and first-match joins.
- Correctness / compatibility analysis: Compared relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Found one introduced P2 regression: with top-K fusion enabled, AQE replanning can fail queries containing entirely valid Base64. Previously reported scan-fallback and AQE assertion concerns are addressed.
- Key design decisions: Early aggregate tagging preserves buffer compatibility. Reusing the existing aggregate-restoration pass avoids another overlapping implementation. Synthetic native operators need separate handling because their saved Spark plans do not necessarily describe their actual evaluation.
- Implementation sketch: Cached serde enrollment drives persistent planner tags, transition removal and restoration of affected native ancestors. Exchanges allow native execution to resume.
- Behavioral changes worth calling out: Compared with
branch-1.1, protected pipelines now fall back by default, including for valid input. The documentedspark.comet.exec.preserveEvaluationMasks.enabled=falseopt-out retains acceleration. Runtime performance cost was not benchmarked. Broader throwing-expression coverage remains tracked in #6006. - Suggested improvements: Address the fused top-K restoration regression and add an AQE execution regression covering its output projection and offset. Request changes for the finding below.
Reviewed full SHA 8401da22890fbdbbb3225b8fae5ececa2c45e9a6: all three commits and 12 files in the PR diff against base ac9ae94d057013095fdf1b70e2d4bc5b24a58d90, using merge base 98662215dc0f485cbbedd4b9a17de8c62fe1fb3f. Confirmed non-draft status. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-shuffle-pr. Read all supplied discussion and threads, excluding Copilot.
Exact-head CI: 51 checks passed and 15 were skipped, with none failed or pending. Run 37222352849 passed Required Checks, all 20 Comet profile jobs and all nine Spark 4.1 SQL shards. Logs confirm the nine evaluation-mask tests and both SQL-fixture dispatcher variants passed.
Local validation: Current-source JVM test compilation, whitespace and CI-registration checks passed. Using the SHA256-verified native artifact for this head, all nine evaluation-mask tests passed. Five additional probes produced two failures demonstrating the same top-K defect. All five passed with the exact base CometExecRule substituted. Disabling fusion or mask protection also passed. Local runtime testing was limited to Spark 4.1.3. Native compilation, macOS and other upstream Spark SQL profiles were not run locally. Project source remains unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Batch execution could decode malformed
unbase64values in rows Spark skips, failing otherwise successful queries. - Design approach:
RequiresSparkEvaluationMaskenrolls expressions in planner protection that preserves Spark row execution around early-stop operators. - Correctness / compatibility analysis: Compared decoder, limit, join and window behavior against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Checked AQE buffer compatibility, repeated planning and fused top-K projection/offset handling. No introduced P1/P2 issues found within this review.
- Key design decisions: Early aggregate tagging protects buffer compatibility before materialization. Reusing the existing aggregate-restoration pass avoids another overlapping implementation. The local top-K selector is distinguished from its enclosing projection and offset.
- Implementation sketch: Cached serde enrollment drives persistent planner tags, removal of intervening batching transitions and restoration of affected native ancestors. Exchanges permit native execution to resume, and pruning-only scans remain native.
- Behavioral changes worth calling out: Compared with
branch-1.1, affected pipelines intentionally fall back by default, including for valid input.spark.comet.exec.preserveEvaluationMasks.enabled=falseprovides the documented opt-out. The early bailout limits unnecessary planning work, but performance costs were not benchmarked. Broader throwing-expression coverage remains tracked in #6006. - Suggested improvements: None meeting the P1/P2 reporting bar. Previously reported scan fallback, AQE assertion and fused top-K concerns are addressed.
Reviewed full SHA c050c9510352afe8a57d018d69f857aa0fe097b2, covering all five commits and 12 PR files against target base ac9ae94d057013095fdf1b70e2d4bc5b24a58d90, using merge base 98662215dc0f485cbbedd4b9a17de8c62fe1fb3f. Confirmed non-draft status. Read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-shuffle-pr.
Exact-head CI: 52 checks passed and 15 were skipped, with none failed or unfinished. Run 37229667890 passed all 20 Linux Comet profile jobs, all nine Spark 4.1 SQL shards and Required Checks. These jobs tested GitHub’s merge commit containing this head and newer main changes. macOS and other upstream Spark SQL profiles were skipped.
Local validation: Current-source JVM compilation, whitespace and CI-registration checks passed. On Spark 4.1.3, all ten evaluation-mask tests and an additional malformed-output top-K probe passed. The probe covers outer limits, offsets, repeated planning, fusion settings and one/two input partitions. Final validation used a SHA256-verified native artifact from base 98662215, whose native source tree exactly matches this head. Native compilation, other runtime profiles and benchmarks were not run locally. Project source remains unchanged.
Which issue does this PR close?
Closes #5532. The broader throwing-expression family is tracked in #6006.
Rationale for this change
Comet can evaluate malformed Base64 in rows that Spark skips below a limit or after a semi/anti join finds its first match. Preserve Spark's evaluation at those operator boundaries so skipped malformed input does not fail the query.
What changes are included in this PR?
The shared
RequiresSparkEvaluationMaskserde policy enrollsunbase64, including strict and dispatched forms. The affected row pipeline stays in Spark below limits, preordered top-K and unpartitioned window limits, and in first-match join conditions. Materialization boundaries can restart native execution. AQE chooses compatible aggregate buffers before materialization and uses main's existing aggregate restoration pass.spark.comet.exec.preserveEvaluationMasks.enableddefaults to true. Users with known-valid input can disable this operator-level protection; expression-level restrictions remain. Plans without enrolled expressions or sticky tags return before the protection walk. This does not claim complete parity for every throwing expression or conditional context; #6006 records the remaining scope.V1 scan predicates only count as row evaluation sites when native Parquet row-filter pushdown is enabled. With the default pruning-only path, the scan stays native while the Spark Filter above it preserves row evaluation. Regression assertions also permit either valid inner-join broadcast side and require both compatible aggregate stages to remain native.
The fused local top-K is treated as a candidate selector that drains its input, without inheriting the enclosing top-K's projection or offset. Repeated planning and AQE therefore apply those final operations once, including when an outer limit requires Spark fallback.
How are these changes tested?
All 10
CometEvaluationMaskSuitetests pass on Spark 4.1.3. The new fused top-K regression covers AQE on/off, fusion on/off, broadcast-stage execution, repeated rule application, nonzero offset and an outer-limit fallback. Before the fix, the same regression reproduces the reported AQE broadcast-stage error while fusion-disabled controls pass.Strict Spark 3.5 / Scala 2.12 main and test compilation and the focused fused top-K runtime regression pass. Spotless, Scala style, suite-registration and whitespace checks pass. Native source/build inputs exactly match base
98662215dc0f485cbbedd4b9a17de8c62fe1fb3f; JVM tests use its verified CI native artifact (run 37219565655).Multi-profile and upstream Spark SQL CI are requested separately; local results do not establish those jobs have passed.