Skip to content

[SPARK-59122][4.3][SQL] Take UnionExec's plain-union decision once instead of re-deriving it per caller - #59151

Open
LuciferYang wants to merge 1 commit into
apache:branch-4.3from
LuciferYang:SPARK-59122-4.3
Open

LuciferYang wants to merge 1 commit into
apache:branch-4.3from
LuciferYang:SPARK-59122-4.3

Conversation

@LuciferYang

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

This backports SPARK-59122 (#58419, already in master and branch-4.x) to branch-4.3. The change makes UnionExec decide once, right after EnsureRequirements, whether it is a plain concatenation, and freeze that decision plus the codegen confs in a TreeNodeTag, so the fusion gate, metrics, unionRDDs and the copy of the union inside the codegen shell can no longer give different answers depending on when each is read. See #58419 for the full design and the two failures it fixes.

Two conflicts came up against branch-4.3, both because the branch predates code the master change was built on. Neither changes the fix:

  • branch-4.3 does not have SPARK-57399, so AQEEnablePipelinedShuffle is not in the AQE query-stage preparation rule list. The two union barriers are appended after RemoveRedundantSorts exactly as on master, without that rule.
  • branch-4.3 does not have SPARK-58819 or SPARK-59285, so UnionExec still derives its candidate partitioning by remapping to the first child's attributes and lifting with a local toUnionOutput, and comparePartitioning does not compare keyDataTypes. The backport keeps that body. The SPARK-59122 edit to it is only the rename of the old outputPartitioning body to a private rawPartitioning, moving the spark.sql.unionOutputPartitioning gate out of an early return and into the stamped decision, and adding the barriers and the preparation-conf snapshot around it.

Why are the changes needed?

SPARK-59122 is a Blocker correctness fix. On default settings a union whose children share a partitioning could concatenate at execution while a parent had already dropped the exchange it needed, giving wrong results; a fused union could also lose the numOutputRows metric its generated code asks for and crash. branch-4.3 is exposed to both, since spark.sql.unionOutputPartitioning defaults to true there.

Does this PR introduce any user-facing change?

Yes, the same change as #58419: the wrong-result and codegen-crash cases above now behave correctly. The three confs the decision reads no longer reach a plan that has already been prepared, only plans prepared afterwards.

How was this patch tested?

The tests come with the cherry-pick, unchanged from #58419; CI runs them against branch-4.3. The AQE rule-order case asserts positions relative to EnsureRequirements and the end of the list, so it holds with AQEEnablePipelinedShuffle absent without editing.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Opus 5

… of re-deriving it per caller

`UnionExec` re-derived its answers from the children's `outputPartitioning` on every call, and that answer moves between planning and execution, so the fusion gate, `metrics` and `unionRDDs` could disagree with each other.

What moves is now decided once per node and kept in a `TreeNodeTag`: whether the union is a plain concatenation (the old `outputPartitioning` body becomes a private `rawPartitioning`) and the two codegen confs. A new `StampUnionDecisions` rule asks for that answer right after `EnsureRequirements` in the standard and the AQE preparation pipeline, and again after each phase that can add a `UnionExec` of its own (the injected columnar and query-stage rules). The confs are recorded one rule earlier, by `SnapshotUnionPreparationConf`, so `EnsureRequirements` and the stamp behind it cannot sample `spark.sql.unionOutputPartitioning` separately. Every barrier in one preparation takes the same `UnionConfSnapshot`, read once, so a union an injected rule created or rebuilt after the first pass is stamped from the values the exchanges above it were planned against and not from the conf as it is by then. Where rules are folded over rather than listed, the injected columnar rules and the AQE query-stage optimizers, the same record is written behind each rule that changed the plan, so a union one of them returns is read from the preparation's values by the next rule, and by `ValidateRequirements` where that rule is itself an `AQEShuffleReadRule`. Reads before the barrier answer from what they see and write nothing, so observing an unprepared plan decides nothing for the prepared one.

A tag rather than a field, because `withNewChildren` ends in `copyTagsFrom`: the copy `CollapseCodegenStages` puts inside the codegen shell inherits the stamped values instead of taking its own. The gate's remaining terms read the children, so they stay per-instance and are re-derived when a rule installs different children.

Two failures on default configuration, both from those reads disagreeing.

A crash, when a fused union loses the metric its generated code needs:

```
java.util.NoSuchElementException: key not found: numOutputRows
  at org.apache.spark.sql.execution.CodegenSupport.metricTerm(WholeStageCodegenExec.scala:71)
  at org.apache.spark.sql.execution.UnionExec.doProduce(basicPhysicalOperators.scala:1148)
  at org.apache.spark.sql.execution.WholeStageCodegenExec.doCodeGen(WholeStageCodegenExec.scala:676)
```

```scala
spark.range(0, 200, 1, 4).selectExpr("id % 10 AS k", "id AS v")
  .groupBy("k").agg(sum("v").as("s")).createOrReplaceTempView("v")
spark.catalog.cacheTable("v")
spark.sql("SELECT k, abs(s) AS s FROM v UNION ALL SELECT k, s FROM v").collect()
```

`InMemoryTableScanExec` reports `UnknownPartitioning` while its inner `AdaptiveSparkPlanExec` has no final plan, so the union looks plain when `CollapseCodegenStages` gates on it and is fused. The copy `insertInputAdapter` puts in the shell evaluates the gate only after the cache stages finalise; by then both children report the same `HashPartitioning`, so `metrics` comes back empty under code that asks for `numOutputRows`. Flipping `spark.sql.codegen.wholeStage.union.enabled` between planning and execution reaches the same crash with no cache involved (given one child that is not `CodegenSupport`), which is why the codegen confs are stamped too. Registering the metric unconditionally stops both crashes, but leaves the wrong answer below in place and puts a 0-valued row count on every union that falls back to `doExecute`.

And a wrong answer, because `rawPartitioning` read `spark.sql.unionOutputPartitioning` live while `executedPlan` is memoized on first read:

```scala
val df = spark.range(0, 20, 1, 2).selectExpr("id % 5 AS k").repartition(4, col("k"))
  .union(spark.range(20, 40, 1, 2).selectExpr("id % 5 AS k").repartition(4, col("k")))
  .groupBy("k").count()

df.explain()   // the union reports HashPartitioning(k, 4), so the aggregate's exchange is elided
spark.conf.set("spark.sql.unionOutputPartitioning", "false")
df.collect()   // the union concatenates instead, so each group lands in two partitions
```

Five groups of eight come back as ten rows of four. A prepared node can no longer change what it reports or how it executes, which fixes it. Deciding on first read would not have: `executedPlan` is `prepareForExecution(sparkPlan.clone())` and `clone` ends in `copyTagsFrom`, so a read on `sparkPlan` would have decided for the prepared plan, and where nothing reads the node during preparation the first read lands at execution anyway.

Half of that gap stays open: a union stamped non-plain still derives `rawPartitioning` per call, so a partitioning change on one child that its siblings do not mirror can leave it concatenating at execution after reporting something concrete at planning. AQE reverts such a change when it breaks a parent's requirement, so what remains is a rule injected at one of the extension points, and a check at execution cannot close that, since the node cannot tell whether a parent relied on what it reported.

This affects the fusion added in SPARK-56482, so branch-4.2 onward carries it.

Yes. The first query above fails on 4.2.0 and now returns rows; the second returned duplicated groups and now returns the correct ones.

The three confs these decisions read (`spark.sql.unionOutputPartitioning`, `spark.sql.codegen.wholeStage.union.enabled`, `spark.sql.codegen.wholeStage.union.maxChildren`) no longer reach a plan that has already been prepared, only plans prepared afterwards. A DataFrame built before a change but first prepared after it observes the new value.

Three more changes come with deciding early, none of them affecting results:

- A union whose children only agree on partitioning after planning keeps reporting `UnknownPartitioning`, so SPARK-52921's exchange elimination no longer applies to that shape. The alternative is a node that concatenates partitions while advertising a partitioning it does not have.
- `DisableUnnecessaryBucketedScan` runs after the first stamping pass, so a union over two bucketed scans with a projection on each side is stamped from the bucketed partitioning and then loses it: the gate answers `partitioning-aware`, and the union does not fuse, whereas deriving at the gate would have fused it. Closing that means splitting the two decisions, which is a separate change.
- `CoalesceShufflePartitions` keeps a non-plain union's children in one coalesce group, so if `OptimizeSkewInRebalancePartitions` splits only one of them, both lose coalescing where each used to be coalesced independently.

In `UnionCodegenSuite`, for the two failures above and the decisions they turn on:

- a fused union keeps `numOutputRows`, and reports `UnknownPartitioning` rather than its children's `HashPartitioning`
- a union planned with `spark.sql.unionOutputPartitioning` on still executes that way after the conf is flipped off
- the same for `spark.sql.codegen.wholeStage.union.enabled`, over exchange children so the shell really holds a copy
- the same for `spark.sql.codegen.wholeStage.union.maxChildren`, over a three-child union whose cap is lowered under it
- an already-decided gate answers against the children a later `withNewChildren` gives it
- a read on the unprepared plan does not decide for the prepared one
- a read before any barrier leaves the gate provisional, so an unstamped union registers `numOutputRows` whichever way the stamp later lands
- a partitioning-aware union follows its children's coalesced partition count, so the count `unionRDDs` hands `SQLPartitioningAwareUnionRDD` is never a stamped one
- a union fused under AQE keeps both codegen decisions through execution, where the shell is built while the query runs

And in the same suite, for the barriers and the snapshot, driving the new rules by hand:

- a later stamping pass fills in a union an extension made, and leaves an earlier decision alone
- the stamp answers from the conf recorded before `EnsureRequirements`, not from the value live when it runs
- a replacement node whose own tag kept `copyTagsFrom` from bringing the stamp across is stamped by a late barrier from the preparation's conf, not from the value live when it runs
- a barrier reaching a union no snapshot pass saw records that preparation's codegen confs rather than the values live when it runs
- a snapshot pass ahead of each injected rule fills in a union the rule before it created
- a recorded conf answers the codegen gate before any decision is stamped

In `SparkSessionExtensionSuite`, through the real pipelines instead of calling the rules:

- a union an injected columnar rule adds is stamped with AQE off, and so is one the columnar hook adds inside AQE post stage creation
- so is one an injected query stage prep rule adds. That case reads the node from an injected stage-optimizer rule, since the barrier in `postStageCreationRules` otherwise stands behind the one after the prep rules and would stamp it anyway
- a union one injected prep rule adds is recorded before the next prep rule reads it, and the same for two injected post planner strategy rules, and again through the codegen gate rather than the partitioning
- a union one injected query stage optimizer rule adds is recorded before the next one reads it, and the same for two injected columnar rules, in each transition phase and with AQE off and on
- with `spark.sql.unionOutputPartitioning` flipped after the `AdaptiveSparkPlanExec` wrapper has read its snapshot and before the stage carrying the union is created, the barrier there still stamps from the snapshot
- the same for the two codegen confs, over a three-child union an injected columnar rule rebuilds: it still fuses, and the copy inside the shell counts its rows
- an injected `AQEShuffleReadRule` puts a union under the final aggregate, and the rewrite survives `ValidateRequirements` only because the union answers from the record rather than from the flipped conf

Each fails when the barrier, the snapshot pass, or the record it covers is removed. The last two also fail when that barrier reads the conf live rather than taking the wrapper's snapshot: the partitioning one reports `UnknownPartitioning(0)`, the codegen one fuses nothing.

One case in `AdaptiveQueryExecSuite` pins the order of the AQE list: the snapshot pass immediately before `EnsureRequirements`, the first stamp immediately after it, and the barrier for the injected prep rules last.

Generated-by: Claude Opus 5

Closes apache#58419 from LuciferYang/SPARK-59122.

Authored-by: YangJie <yangjie01@baidu.com>
Signed-off-by: yangjie01 <yangjie01@baidu.com>
(cherry picked from commit 8f6d613)
@LuciferYang

Copy link
Copy Markdown
Contributor Author

cc @cloud-fan

@uros-b
uros-b requested a review from cloud-fan September 30, 2026 11:13
@uros-b

uros-b commented Sep 30, 2026

Copy link
Copy Markdown
Member

Thank you @LuciferYang!

This branch has not been deployed

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants