cloud-fan commented on code in PR #58419:
URL: https://github.com/apache/spark/pull/58419#discussion_r4036651738
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/basicPhysicalOperators.scala:
##########
@@ -1284,6 +1391,35 @@ case class UnionExec(children: Seq[SparkPlan]) extends
SparkPlan with CodegenSup
}
object UnionExec {
+ /**
+ * What `StampUnionDecisions` fixes on a `UnionExec`: whether it is a plain
concatenation, and the
+ * two confs the codegen gate reads. Everything else the gate asks is
derived per instance, so a
+ * rule replacing a child cannot inherit an answer taken from the topology
it replaced.
+ */
+ private case class Decisions(
+ plainUnion: Boolean,
+ unionCodegenEnabled: Boolean,
+ maxChildren: Int)
+
+ /**
+ * The stamped decisions. See `isPlainUnion` and `stampDecisions`.
+ *
+ * `withNewChildren` copies the tag onto a rebuilt node, and so does a
transform rule's
+ * replacement, but only where the target carries no tags of its own:
`copyTagsFrom` leaves a node
+ * that already has some untouched. A `UnionExec` reaching execution
unstamped therefore answers
+ * from the state it sees then, and can leave `metrics` empty, so
`doProduce` fails asking
+ * `metricTerm` for `numOutputRows`.
+ */
+ private val DECISIONS = TreeNodeTag[Decisions]("unionDecisions")
Review Comment:
**Non-blocking (P2):** `copyTagsFrom` is all-or-nothing: if a replacement
already carries any tag, neither `OUTPUT_PARTITIONING_CONF` nor `DECISIONS` is
copied. These values now bridge `EnsureRequirements` and execution, so a
supported late extension can return an equivalent `UnionExec` with unrelated
metadata, lose both values, and be restamped from the live configuration. If
the setting changes in that interval, a parent may already have omitted its
exchange based on the old concrete partitioning while the replacement executes
as a plain concatenation, yielding duplicate aggregate groups. Please keep this
preparation state explicitly on `UnionExec` outside its primary case-class
product, forward it through intentional copies, and make late barriers use the
same preparation-scoped configuration snapshot rather than silently falling
back to the live conf.
**Recommended change:** Replace OUTPUT_PARTITIONING_CONF and DECISIONS
TreeNodeTag values with explicit non-product UnionExec preparation state, and
give the preparation rules one immutable configuration snapshot shared by the
initial and late barriers. Preserve existing-node state through intentional
UnionExec copy paths, prepare genuinely new extension nodes from that shared
snapshot, and fail closed if an unprepared union reaches codegen or execution.
**Why this works:** Represent unprepared, configuration-captured, and
prepared states in private UnionExec state outside the primary case-class
parameter list, forwarding it deliberately through withNewChildrenInternal and
clone-sensitive copies. Capture the union configuration once in a
preparation-scoped immutable input shared by the rules before and after
EnsureRequirements and by the late extension barriers, so a genuinely new
extension union is prepared from the same input rather than the live conf. Add
a boundary check so missing preparation context cannot fall back to a new live
configuration during codegen or execution.
**Scope:** Make UnionExec decision transport explicit and fail closed when
preparation state is missing.
**Compatibility:** The preparation-time configuration lifecycle, write-once
decisions, InputAdapter copy agreement, and dynamic AQE partition counts remain
unchanged for correctly prepared plans.
**Risks:** A direct UnionExec copy site omitted from explicit state
forwarding could fail validation or regress the original mismatch. Including
preparation state in structural plan equality or canonicalization could inhibit
reuse and destabilize plan output.
**Constraints:** Unprepared-plan inspection remains non-mutating and may
answer from live configuration. The output-partitioning configuration sampled
before EnsureRequirements remains the one used for the corresponding completed
decision. Only the plain-versus-aware decision is retained; the Partitioning
object and AQE partition count remain derived per call. Plan canonicalization,
reuse identity, and normal explain output must not acquire preparation-state
semantics.
**Success:** An unrelated TreeNode tag on a replacement cannot erase
UnionExec's configuration snapshot or completed decision. Every copy of an
existing prepared UnionExec used by requirements, codegen, metrics, or
execution observes the same completed decision. A genuinely new
extension-created UnionExec is prepared once at its supported barrier. No
unprepared UnionExec can silently reach codegen or execution and sample live
configuration.
##########
sql/core/src/test/scala/org/apache/spark/sql/execution/UnionCodegenSuite.scala:
##########
@@ -628,6 +694,346 @@ class UnionCodegenSuite extends SharedSparkSession {
}
}
+ test("SPARK-59122: a fused union keeps numOutputRows and reports
UnknownPartitioning") {
+ // The children's partitioning is not stable while the plan is being
prepared:
+ // `InMemoryTableScanExec.cachedPlan` unwraps the inner
`AdaptiveSparkPlanExec` only once
+ // `isFinalPlan` is true, and reports `UnknownPartitioning(0)` until then,
so the union looks
+ // plain and is fused. The projection is what makes that reachable:
`supportsColumnar` is
+ // `children.forall`, so one row-based `ProjectExec` over the columnar
scan is enough to make
+ // it false, and without one `supportCodegenFailureReason` reports
`columnar` and nothing
+ // fuses. `SELECT *` or a plain alias collapses the projection away and
does not reproduce
+ // this. Once the cache stages finalise, both children report the same
concrete layout, and
+ // re-deriving the decision at that point left `metrics` empty while
`doProduce` asked
+ // `metricTerm` for `numOutputRows`.
+ //
+ // Both halves of the decision are asserted here. Registering the metric
unconditionally would
+ // fix the crash and leave the other half broken: a fused union
concatenates its children's
+ // partitions, so claiming their partitioning would let a parent satisfy a
clustered
+ // distribution from an RDD that does not have it.
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true",
+ SQLConf.WHOLESTAGE_UNION_CODEGEN_ENABLED.key -> "true",
+ SQLConf.UNION_OUTPUT_PARTITIONING.key -> "true") {
+ withTempView("v") {
+ cacheAggregateView("v")
+ val df = spark.sql("SELECT k, abs(s) AS s FROM v UNION ALL SELECT k, s
FROM v")
+ // Execute this DataFrame rather than a count over it: the plan being
inspected has to be
+ // the one that ran, and an AQE plan that never ran has no final plan
to inspect.
+ assert(df.collect().length == 20)
+ val fused = fusedUnions(df)
+ assert(fused.nonEmpty,
+ "this shape must actually fuse, or the test is not exercising the
defect")
+ fused.foreach { u =>
+ // Part of the premise, not the whole of it: the children expose a
concrete layout by now,
+ // so this node is not reporting `UnknownPartitioning` merely for
want of anything to
+ // derive from. `rawPartitioning` also falls back when the
children's remapped
+ // partitionings do not compare equal, and that cannot be asserted
here: each side carries
+ // its own exprIds, and they line up only after the private
`prepareOutputPartitioning`.
+ val childPartitionings = u.children.map(_.outputPartitioning)
+
assert(childPartitionings.forall(_.isInstanceOf[HashPartitioningLike]),
+ s"premise: got $childPartitionings")
+ assert(childPartitionings.map(_.numPartitions).distinct.size == 1,
+ s"premise: got $childPartitionings")
+ assert(u.metrics.contains("numOutputRows"),
+ "a fused union must register the metric its generated code
increments")
+ assert(u.outputPartitioning.isInstanceOf[UnknownPartitioning],
+ s"a fused union must not claim a concrete partitioning, got
${u.outputPartitioning}")
+ }
+ }
+ }
+ }
+
+ test("SPARK-59122: a partitioning-aware union keeps its layout when the conf
changes between " +
+ "planning and execution") {
+ // `spark.sql.unionOutputPartitioning` is read once during preparation,
ahead of
+ // `EnsureRequirements`, not on every `outputPartitioning` call, so a plan
executes by the
+ // partitioning it was planned against. Reading it per call let the parent
aggregate lose its
+ // exchange at planning and get a plain concatenation at execution,
reporting each group twice.
+ // The `checkAnswer` below stays outside the block that planned the
DataFrame on purpose: the
+ // plan is forced inside that block and `executedPlan` is memoized, so the
two phases see
+ // different confs. Asserting inside it, or dropping the second
`withSQLConf`, makes the test
+ // pass without testing this.
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ val left = spark.range(0, 20, 1, 2).selectExpr("id % 5 AS k")
+ val right = spark.range(20, 40, 1, 2).selectExpr("id % 5 AS k")
+
+ val planned = withSQLConf(SQLConf.UNION_OUTPUT_PARTITIONING.key ->
"true") {
+ val df = left.repartition(4, col("k"))
+ .union(right.repartition(4, col("k"))).groupBy("k").count()
+ val plan = df.queryExecution.executedPlan
+ val unions = plan.collect { case u: UnionExec => u }
+ assert(unions.size == 1)
+ // Asserted through the exchanges rather than through `isPlainUnion`,
so that the check
+ // does not depend on how the decision is stored: only the two
repartitions may shuffle, so
+ // the aggregate's exchange was elided, which it could only be if the
union reported a
+ // concrete partitioning.
+ val shuffles = plan.collect { case s: ShuffleExchangeExec => s }
+ assert(shuffles.size == 2)
+ assert(shuffles.forall(_.shuffleOrigin == REPARTITION_BY_NUM),
+ s"expected only the two repartitions, got
${shuffles.map(_.shuffleOrigin)}")
+ df
+ }
+
+ withSQLConf(SQLConf.UNION_OUTPUT_PARTITIONING.key -> "false") {
+ // Each side contributes four ids per `k`, so the answer is fixed.
Comparing against the
+ // same query run with the conf off would also pass if both paths
regressed to ten rows.
+ checkAnswer(planned, (0L until 5L).map(k => Row(k, 8L)))
+ }
+ }
+ }
+
+ test("SPARK-59122: a fused union keeps numOutputRows when the codegen conf
changes between " +
+ "planning and execution") {
+ // `supportCodegenFailureReason` used to read
`WHOLESTAGE_UNION_CODEGEN_ENABLED` live, and the
+ // copy that `insertInputAdapter` puts inside the codegen shell evaluated
it for the first time
+ // at execution. Planned with the conf on the union is fused, so the
generated code increments
+ // `numOutputRows`; if the copy re-derives the reason with the conf off,
`metrics` comes back
+ // empty and `doProduce` throws `key not found: numOutputRows`.
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ val planned = withSQLConf(SQLConf.WHOLESTAGE_UNION_CODEGEN_ENABLED.key
-> "true") {
+ // Each child is an exchange, which is not `CodegenSupport`, so
`insertInputAdapter` wraps
+ // it and the union is rebuilt through `withNewChildren`, the copy
this test needs. Children
+ // that do support codegen can still produce one, since
`insertInputAdapter` recurses into
+ // their descendants; exchanges just make it certain.
+ val df = rangeDF(100).repartition(2).union(rangeDF(100).repartition(2))
+ // `fusedUnions` requires the union to be the stage root, which is
what this test needs: it
+ // reaches for that node itself below.
+ assert(fusedUnions(df).size == 1, "this shape must fuse, or the test
exercises nothing")
+ df
+ }
+ withSQLConf(SQLConf.WHOLESTAGE_UNION_CODEGEN_ENABLED.key -> "false") {
+ assert(planned.collect().length == 200)
+ // The row count alone does not discriminate, since the shell was
installed at planning and
+ // keeps emitting; registering `numOutputRows` unconditionally and
reading the conf per call
+ // passes it. The `supportCodegen` assertion below is what fails there.
+ val copy = fusedUnions(planned)
+ assert(copy.size == 1)
+ // The copy this test needs: `insertInputAdapter` wrapped both
children, so the shell holds
+ // a copy rather than the instance the gate answered on. This copy's
reason is first forced
+ // by the `SparkPlanInfo` that `collect()` above builds, with the conf
already off, so what
+ // it answers can only come from the stamp.
+ assert(copy.head.children.forall(_.isInstanceOf[InputAdapter]))
+ assert(copy.head.supportCodegen,
+ "the copy in the shell must keep the decision it was planned with")
+ }
+ }
+ }
+
+ test("SPARK-59122: a fused union keeps numOutputRows when the child cap
drops between " +
+ "planning and execution") {
+ // `WHOLESTAGE_UNION_MAX_CHILDREN` is on the same snapshot as the enable
flag, so the same shape
+ // has to hold for it: prepared under a cap this union meets, it stays
fused even if the cap is
+ // lowered under it. Reading the cap live would give the shell's copy
`max-children-exceeded`,
+ // empty `metrics`, and `doProduce` failing at `metricTerm`. Three
children against a cap of
+ // two, since the conf refuses anything below two.
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ val planned = withSQLConf(
+ SQLConf.WHOLESTAGE_UNION_CODEGEN_ENABLED.key -> "true",
+ SQLConf.WHOLESTAGE_UNION_MAX_CHILDREN.key -> "3") {
+ // Exchange children again, so the shell really holds a
`withNewChildren` copy.
+ val df = rangeDF(100).repartition(2)
+ .union(rangeDF(100).repartition(2))
+ .union(rangeDF(100).repartition(2))
+ val fused = fusedUnions(df)
+ assert(fused.size == 1 && fused.head.children.size == 3,
+ s"this shape must fuse as one three-child union, got
${fused.map(_.children.size)}")
+ df
+ }
+ withSQLConf(SQLConf.WHOLESTAGE_UNION_MAX_CHILDREN.key -> "2") {
+ assert(planned.collect().length == 300)
+ val copy = fusedUnions(planned)
+ assert(copy.size == 1)
+ assert(copy.head.children.forall(_.isInstanceOf[InputAdapter]))
+ assert(copy.head.supportCodegen,
+ "the copy in the shell must keep the cap it was planned with")
+ // Not `metrics.contains`, which `collect()` above already proves: an
empty `metrics` would
+ // have thrown at `metricTerm`. The count is what says the fused code
ran and counted.
+ assert(copy.head.metrics("numOutputRows").value == 300)
+ }
+ }
+ }
+
+ test("SPARK-59122: the codegen gate re-derives when a rule replaces the
children") {
+ // The gate's children-dependent terms must not outlive the children they
were taken from.
+ // `SQLExecution` builds a `SparkPlanInfo` before execution, which reads
`metrics` on every
+ // node; a decision carried from there onto a node whose children a rule
then replaced would
+ // fuse a topology that the gate rejects.
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.WHOLESTAGE_UNION_CODEGEN_ENABLED.key -> "true") {
+ val df = rangeDF(100).union(rangeDF(100))
+ val unions = fusedUnions(df)
+ assert(unions.size == 1, "this shape must fuse, or the test exercises
nothing")
+ val union = unions.head
+ // What the plan update does, and what decides the gate for this
instance.
+ assert(union.metrics.contains("numOutputRows"))
+ assert(union.supportCodegen)
+
+ // A nested union is one of the topologies the gate rejects, and
`withNewChildren` is the path
+ // a rule takes when it rewrites children in place. A rule returning an
arbitrary replacement
+ // node is a different path, and one `copyTagsFrom` need not carry the
tags along.
+ val nested = UnionExec(Seq(union.children.head, union.children.head))
+ val rebuilt = union.withNewChildren(Seq(nested,
union.children.last)).asInstanceOf[UnionExec]
+ assert(!rebuilt.supportCodegen, "the rebuilt union must answer against
its own children")
+ // Implied by the line above as the code stands, and kept as the pin on
that: registering the
+ // metric unconditionally would leave the line above green, and only
this one would fail.
+ assert(rebuilt.metrics.isEmpty)
+ }
+ }
+
+ test("SPARK-59122: reading the unprepared plan does not decide the prepared
one") {
+ // `QueryExecution.executedPlan` is
`prepareForExecution(sparkPlan.clone())`, and `clone` ends
+ // in `makeCopy`, which calls `copyTagsFrom`. A decision written while
answering a read on
+ // `sparkPlan` would therefore ride into the prepared plan. Here the two
answers differ: each
+ // child is an aggregate whose exchange `EnsureRequirements` has yet to
insert, so the union
+ // passes nothing through before preparation and both children's
`HashPartitioning` after it.
+ // Reads before `StampUnionDecisions` answer without writing, so only
preparation decides.
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.UNION_OUTPUT_PARTITIONING.key -> "true") {
+ val left = spark.range(0, 20, 1, 2).selectExpr("id % 5 AS
k").groupBy("k").count()
+ val right = spark.range(20, 40, 1, 2).selectExpr("id % 5 AS
k").groupBy("k").count()
+ val df = left.union(right)
+
+ val unprepared = df.queryExecution.sparkPlan.collect { case u: UnionExec
=> u }
+ assert(unprepared.size == 1)
+
assert(unprepared.head.outputPartitioning.isInstanceOf[UnknownPartitioning],
+ "the aggregates have no exchange under them yet, so there is nothing
to pass through")
+
+ val prepared = df.queryExecution.executedPlan.collect { case u:
UnionExec => u }
+ assert(prepared.size == 1)
+
assert(!prepared.head.outputPartitioning.isInstanceOf[UnknownPartitioning],
+ "the read above must not have decided for the prepared plan, got " +
+ s"${prepared.head.outputPartitioning}")
+ checkAnswer(df, (0L until 5L).flatMap(k => Seq(Row(k, 4L), Row(k, 4L))))
+ }
+ }
+
+ test("SPARK-59122: a partitioning-aware union follows its children's
coalesced partition count") {
+ // Only the decision is stamped, never the `Partitioning`. AQE coalescing
changes the children's
+ // `numPartitions` after the stamp, and `unionRDDs` hands whatever it
reports to
+ // `SQLPartitioningAwareUnionRDD`, which builds exactly that many
partitions from each child: a
+ // count frozen at stamping time asks for partitions the coalesced
children no longer have.
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "true",
+ SQLConf.COALESCE_PARTITIONS_ENABLED.key -> "true",
+ SQLConf.SHUFFLE_PARTITIONS.key -> "20",
+ SQLConf.UNION_OUTPUT_PARTITIONING.key -> "true") {
+ val left = spark.range(0, 100, 1, 4).selectExpr("id % 10 AS
k").groupBy("k").count()
+ val right = spark.range(100, 200, 1, 4).selectExpr("id % 10 AS
k").groupBy("k").count()
+ val df = left.union(right).groupBy("k").agg(sum("count").as("c"))
+ checkAnswer(df, (0L until 10L).map(k => Row(k, 20L)))
+
+ val unions = collect(df.queryExecution.executedPlan) { case u: UnionExec
=> u }
+ assert(unions.size == 1)
+ val children =
unions.head.children.map(_.outputPartitioning.numPartitions)
+ assert(children.distinct.size == 1 && children.head < 20,
+ s"the children must have been coalesced as one group, got $children")
+ assert(unions.head.outputPartitioning.numPartitions == children.head,
+ "the union must report what its children report now, got " +
+ s"${unions.head.outputPartitioning}")
+ }
+ }
+
+ test("SPARK-59122: a later stamping pass fills in a fresh union and keeps
stamped ones") {
+ // `StampUnionDecisions` is listed again after the phases that can add a
`UnionExec`, so one an
+ // injected columnar or query-stage rule created does not answer from
whatever the conf says
+ // wherever it is first asked. A later pass must also not move a decision
already taken, which
+ // is the second half here. The rule is driven directly, since what this
case is about is its
+ // contract; that the pipelines still list it after each phase that can
add a union is pinned
+ // from the outside by the extension-driven cases in
`SparkSessionExtensionSuite`, which need a
+ // session of their own.
+ withSQLConf(SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
+ // Pins the property the standard pipeline has to keep: a stamping pass
runs after
+ // `EnsureRequirements`, so the decision is taken from the plan the
exchanges were placed in.
+ // A count would break on a sixth legitimate pass and say nothing about
the order. The AQE
+ // lists are private to `AdaptiveSparkPlanExec`, so their first pass has
no counterpart here.
+ val rules = QueryExecution.preparations(spark, subquery = false)
+ val firstStamp = rules.indexWhere(_ eq StampUnionDecisions)
+ val ensureRequirements =
rules.indexWhere(_.isInstanceOf[EnsureRequirements])
+ val columnarRules =
+
rules.indexWhere(_.isInstanceOf[ApplyColumnarRulesAndInsertTransitions])
+ assert(ensureRequirements >= 0 && firstStamp > ensureRequirements &&
Review Comment:
**Non-blocking (P2):** This assertion still passes if a
partitioning-changing rule is inserted between `EnsureRequirements` and the
first stamp. That is the unsafe window this lifecycle change needs to close:
the parent can plan against one answer and the stamp can freeze another. Please
assert `firstStamp == ensureRequirements + 1` (while retaining the separate
later-stamp check) so the focused test fails on that regression.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/basicPhysicalOperators.scala:
##########
@@ -1044,11 +1091,21 @@ case class UnionExec(children: Seq[SparkPlan]) extends
SparkPlan with CodegenSup
}
}
- // Memoized: consulted by `supportCodegen` (called multiple times by
- // `CollapseCodegenStages`) and by `metrics`. Conf and children are stable
- // for a given UnionExec instance; cross-plan staleness is impossible since
- // UnionExec is a case class and `withNewChildren` produces a fresh instance.
- @transient private lazy val supportCodegenFailureReason: Option[String] = {
+ // Latched for the same reason `isPlainUnion` is: `supportCodegen` and
`metrics` must see one
+ // answer, and `conf` is live. When a child is not `CodegenSupport`,
`insertInputAdapter` wraps
+ // it, so `withNewChildren` returns a real copy whose first evaluation of
this would land at
+ // execution; re-deriving there left `metrics` empty while `doProduce` asked
`metricTerm` for
+ // `numOutputRows`. The first force is not always the gate: under AQE it is
a plan-update event
+ // on the pre-stage-creation tree, so a term added here sees more of the
plan than the gate does.
+ private def supportCodegenFailureReason: Option[String] =
decisionLock.synchronized {
+ getTagValue(UnionExec.CODEGEN_FAILURE_REASON).getOrElse {
Review Comment:
Confirmed: the query-stage-preparation test now observes the injected union
from a query-stage optimizer before the later post-stage barrier, so removing
the early post-prep stamp fails the case. Resolved.
<!-- SPARK_DEV_REVIEW_REPLY
{"feedback_id":"inline:3979092641","thread_id":"inline:3979092641","verdict_sha256":"6a099963095a21eb62dcca5d834d1640efeaefed2ee53a50abe468171a2023f3"}
-->
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/basicPhysicalOperators.scala:
##########
@@ -1035,13 +1040,96 @@ case class UnionExec(children: Seq[SparkPlan]) extends
SparkPlan with CodegenSup
}
}
- // True when the codegen path applies: `outputPartitioning` is
`UnknownPartitioning`,
- // and `unionedInputRDD` matches the semantics of `sparkContext.union(...)`
in `unionRDDs`.
- // A `KeyedPartitioning` union also uses `sparkContext.union(...)` in
`unionRDDs`, but
- // codegen is disabled for it (`supportCodegenFailureReason` reports
"partitioning-aware"):
- // the per-partition key descriptor is consumed by a downstream
`GroupPartitionsExec`, and
- // keeping these unions out of whole-stage codegen matches the
`HashPartitioning` union case.
- private[sql] def isPlainUnion: Boolean =
outputPartitioning.isInstanceOf[UnknownPartitioning]
+ /**
+ * True when this union behaves as a plain concatenation, so
`unionedInputRDD` matches the
+ * semantics of `sparkContext.union(...)` in `unionRDDs`. It satisfies the
partitioning gate on
+ * the codegen path, not the whole of it: `supportCodegenFailureReason`
still applies its other
+ * checks. When this union merges its children's `KeyedPartitioning`
instead, it concatenates all
+ * the same, but codegen stays off, with the reason "partitioning-aware",
because a downstream
+ * `GroupPartitionsExec` consumes its key descriptor.
+ *
+ * Stamped, because the answer moves under its consumers.
+ * `InMemoryTableScanExec.outputPartitioning` reports `UnknownPartitioning`
while its inner
+ * `AdaptiveSparkPlanExec` has no final plan, so a union can look plain when
+ * `CollapseCodegenStages` gates on it and partitioning-aware by the time
the stage runs. The
+ * shell that gate builds wraps a `withNewChildren` copy where a child had
to be adapted, and a
+ * copy that re-derived here came back with empty `metrics` while
`doProduce` asked `metricTerm`
+ * for `numOutputRows`. A fresh copy inherits the answer instead, since
`withNewChildren` ends in
+ * `copyTagsFrom`.
+ *
+ * `UNION_OUTPUT_PARTITIONING` is taken from
`snapshotOutputPartitioningConf`, recorded before
+ * `EnsureRequirements`, so the value the exchanges are planned against is
the value execution
+ * uses; a node created after that pass carries no record and reads the live
conf. Reading it live
+ * here would leave one rule between the two: `conf` is live, and another
thread setting it in
+ * that window would let a parent drop an exchange over a concrete
partitioning and then have the
+ * stamp freeze plain concatenation under it.
+ *
+ * A read before `StampUnionDecisions` answers from the children as they are
then, and does not
+ * write, so observing an unprepared plan cannot decide anything for the
prepared one.
+ */
+ private[execution] def isPlainUnion: Boolean =
stampedDecisions.map(_.plainUnion).getOrElse {
+ !outputPartitioningEnabled ||
rawPartitioning.isInstanceOf[UnknownPartitioning]
+ }
+
+ private def stampedDecisions: Option[UnionExec.Decisions] =
+ getTagValue(UnionExec.DECISIONS)
+
+ private def outputPartitioningEnabled: Boolean =
+ getTagValue(UnionExec.OUTPUT_PARTITIONING_CONF)
+ .getOrElse(conf.getConf(SQLConf.UNION_OUTPUT_PARTITIONING))
+
+ /**
+ * Records the conf `isPlainUnion` answers from, read once for the whole
plan by
+ * `SnapshotUnionOutputPartitioningConf` and passed in here, ahead of
`EnsureRequirements`, whose
+ * reads the following stamp has to agree with. Only the conf, never a
partitioning: the exchanges
+ * `EnsureRequirements` adds are not there yet, so a decision taken here
would freeze plain on a
+ * union whose children only become co-partitioned there.
+ */
+ private[execution] def snapshotOutputPartitioningConf(enabled: Boolean):
Unit =
+ if (getTagValue(UnionExec.OUTPUT_PARTITIONING_CONF).isEmpty) {
+ setTagValue(UnionExec.OUTPUT_PARTITIONING_CONF, enabled)
+ }
+
+ /**
+ * Fixes this node's decisions for the rest of the plan's life. Called by
`StampUnionDecisions`,
+ * first right after `EnsureRequirements`, so what the exchanges around this
union were planned
+ * against is what execution uses; the two confs come from one read per plan
there. Nothing else
+ * writes this tag on an existing node, and the nodes the rule writes are
freshly planned and not
+ * yet published, so no reader can be looking at one; `metrics` and the
codegen gate read it
+ * later, and a node that already carries it keeps it, which is how the copy
in the codegen shell
+ * stays in step with the gate.
+ */
+ private[execution] def stampDecisions(codegenEnabled: Boolean, maxChildren:
Int): Unit =
+ if (stampedDecisions.isEmpty) {
+ setTagValue(UnionExec.DECISIONS, UnionExec.Decisions(
+ plainUnion = isPlainUnion,
+ unionCodegenEnabled = codegenEnabled,
+ maxChildren = maxChildren))
+ }
+
+ /**
+ * A node stamped plain reports `UnknownPartitioning` even once its children
agree on a concrete
+ * one: a fused union concatenates, and claiming their partitioning would
let a parent skip an
+ * exchange it needs. The cost is SPARK-52921's exchange elimination for
such a union.
+ *
+ * Only the decision is stamped, never the `Partitioning` itself. AQE
coalescing changes the
+ * children's `numPartitions` after the stamp, and a stale count is what
`unionRDDs` would hand
+ * `SQLPartitioningAwareUnionRDD`, which builds exactly that many partitions
from each child.
+ *
+ * The reverse costs fusion. A rule that runs after the stamp and drops a
child's partitioning
+ * leaves the node stamped non-plain, so the codegen gate answers
"partitioning-aware" and
+ * `numOutputRows` goes unregistered, whereas re-deriving at the gate would
have fused it.
+ * `DisableUnnecessaryBucketedScan` does that to a union over two bucketed
scans with a projection
+ * on each side. Results are unaffected, since the branch below re-derives
and concatenates.
+ *
+ * That branch is derived per call, so `unionRDDs` can take the
concatenating arm even though
+ * `EnsureRequirements` planned the parent against a concrete partitioning:
`comparePartitioning`
+ * compares `HashPartitioningLike` by equality, so a change to one child's
partitioning that its
+ * siblings do not mirror can empty the intersection. AQE reconciles that,
by validating a
+ * partitioning change against the parents' requirements; an injected rule
can skip it.
+ */
+ override def outputPartitioning: Partitioning =
+ if (isPlainUnion) super.outputPartitioning else rawPartitioning
Review Comment:
**Non-blocking (P2):** On an enabled, unstamped node, `isPlainUnion` has
already computed `rawPartitioning` to determine whether it is unknown; when the
result is concrete, this arm computes it again. That repeats the per-child
`AttributeMap` construction and candidate intersections during planning. Please
retain the first raw result within this invocation and reuse it for the return
value, without caching it across calls because AQE partition counts must remain
dynamic.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/QueryExecution.scala:
##########
@@ -852,6 +859,10 @@ object QueryExecution {
RemoveRedundantSorts,
ApplyColumnarRulesAndInsertTransitions(
sparkSession.sessionState.columnarRules, outputsColumnar = false),
+ // A barrier for a `UnionExec` an injected columnar rule just created,
which has no decision
+ // yet and would otherwise take one wherever it is first asked. A
decision already stamped on
+ // a node is kept.
+ StampUnionDecisions,
Review Comment:
**Non-blocking (P2):** This late barrier, and the two analogous AQE suffix
barriers, still traverse the visible plan when no custom rule can introduce a
fresh `UnionExec`. The built-in rules at these positions retain or copy
already-stamped nodes, so the default path pays three no-op tree walks
(including per-stage AQE walks). Please add each suffix stamp only when its
corresponding extension list is non-empty; the post-stage case needs to account
for both injected optimizer and columnar rules. The mandatory initial
post-`EnsureRequirements` stamp should remain unconditional.
**Recommended change:** Append each trailing StampUnionDecisions only when
the immediately preceding extension surface can run a custom rule, while
preserving the mandatory initial stamp and extension-enabled integration
coverage.
**Why this works:** Build the classic columnar suffix from whether columnar
rules are configured; build the AQE prep suffix from injected
query-stage-preparation rules; and build the AQE post-stage suffix from
injected optimizer or columnar rules. Update rule-order coverage so it asserts
the barrier under configured extensions and its absence in default
empty-extension plans.
**Scope:** Make late union-decision barriers conditional on the extension
phases they protect.
**Compatibility:** Prepared unions and extension-created unions retain the
same write-once decisions and execution layout as the pinned implementation.
**Risks:** Omitting a barrier when either an optimizer or columnar extension
can create a union would restore live decisions. Tests that inspect default
rule lists must not mistake an intentionally absent optional barrier for an
ordering regression.
**Constraints:** Always retain the initial snapshot/EnsureRequirements/stamp
sequence. Retain write-once behavior for already stamped unions. A configured
extension-created union must be stamped before the next consumer identified by
the existing integration tests.
**Success:** No late stamping traversal runs for an extension phase whose
relevant custom rule lists are empty. Every supported configured extension
phase that can create UnionExec retains a post-phase stamp. Each
extension-enabled barrier keeps an independent regression signal.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/basicPhysicalOperators.scala:
##########
@@ -1035,13 +1039,79 @@ case class UnionExec(children: Seq[SparkPlan]) extends
SparkPlan with CodegenSup
}
}
- // True when the codegen path applies: `outputPartitioning` is
`UnknownPartitioning`,
- // and `unionedInputRDD` matches the semantics of `sparkContext.union(...)`
in `unionRDDs`.
- // A `KeyedPartitioning` union also uses `sparkContext.union(...)` in
`unionRDDs`, but
- // codegen is disabled for it (`supportCodegenFailureReason` reports
"partitioning-aware"):
- // the per-partition key descriptor is consumed by a downstream
`GroupPartitionsExec`, and
- // keeping these unions out of whole-stage codegen matches the
`HashPartitioning` union case.
- private[sql] def isPlainUnion: Boolean =
outputPartitioning.isInstanceOf[UnknownPartitioning]
+ /**
+ * True when this union behaves as a plain concatenation, so
`unionedInputRDD` matches the
+ * semantics of `sparkContext.union(...)` in `unionRDDs`. It satisfies the
partitioning gate on
+ * the codegen path, not the whole of it: `supportCodegenFailureReason`
still applies its other
+ * checks. When this union merges its children's `KeyedPartitioning`
instead, it concatenates all
+ * the same, but codegen stays off, with the reason "partitioning-aware",
because a downstream
+ * `GroupPartitionsExec` consumes its key descriptor.
+ *
+ * Stamped, because the answer moves under its consumers.
+ * `InMemoryTableScanExec.outputPartitioning` reports `UnknownPartitioning`
while its inner
+ * `AdaptiveSparkPlanExec` has no final plan, so a union can look plain when
+ * `CollapseCodegenStages` gates on it and partitioning-aware by the time
the stage runs. The
+ * shell that gate builds wraps a `withNewChildren` copy where a child had
to be adapted, and a
+ * copy that re-derived here came back with empty `metrics` while
`doProduce` asked `metricTerm`
+ * for `numOutputRows`. A fresh copy inherits the answer instead, since
`withNewChildren` ends in
+ * `copyTagsFrom`.
+ *
+ * `UNION_OUTPUT_PARTITIONING` is read where the decision is stamped rather
than in
+ * `rawPartitioning`, so it too is fixed once the plan is prepared: `conf`
is live, and a plan
+ * must execute by the partitioning it was planned against.
+ *
+ * A read before `StampUnionDecisions` answers from the children as they are
then, and does not
+ * write, so observing an unprepared plan cannot decide anything for the
prepared one.
+ */
+ private[execution] def isPlainUnion: Boolean =
stampedDecisions.map(_.plainUnion).getOrElse {
+ !conf.getConf(SQLConf.UNION_OUTPUT_PARTITIONING) ||
Review Comment:
Confirmed: both pipelines now snapshot the union-output-partitioning setting
before EnsureRequirements and the later stamp uses that recorded value; the
focused flip regression distinguishes a live re-read. Resolved.
<!-- SPARK_DEV_REVIEW_REPLY
{"feedback_id":"inline:4000488270","thread_id":"inline:4000488270","verdict_sha256":"6a099963095a21eb62dcca5d834d1640efeaefed2ee53a50abe468171a2023f3"}
-->
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/adaptive/AdaptiveSparkPlanExec.scala:
##########
@@ -134,7 +134,15 @@ case class AdaptiveSparkPlanExec(
Seq(
CoalesceBucketsInJoin,
RemoveRedundantProjects,
+ // Must run before `ensureRequirements`, which asks a `UnionExec` what
it reports: it
+ // records the conf that answer depends on, so the following
`StampUnionDecisions` freezes the
+ // decision under the same value the exchanges were planned against.
+ SnapshotUnionOutputPartitioningConf,
Review Comment:
Confirmed: AdaptiveQueryExecSuite now asserts the AQE snapshot immediately
before EnsureRequirements, the first stamp immediately after it, and the
trailing extension barrier last. Resolved.
<!-- SPARK_DEV_REVIEW_REPLY
{"feedback_id":"inline:4023166749","thread_id":"inline:4023166749","verdict_sha256":"6a099963095a21eb62dcca5d834d1640efeaefed2ee53a50abe468171a2023f3"}
-->
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala:
##########
@@ -2938,7 +2938,9 @@ object SQLConf {
.doc("When both this conf and `spark.sql.codegen.wholeStage` are true, "
+
"UnionExec participates in whole-stage codegen on its " +
"non-partitioning-aware path: the parent and all children fuse into " +
- "a single WholeStageCodegenExec stage.")
+ "a single WholeStageCodegenExec stage. The value is read when a
UnionExec's " +
+ "decision is fixed during physical preparation, so a change does not
reach a " +
Review Comment:
Confirmed: assertFlagParity now builds a fresh DataFrame under each flag
value and proves the disabled plan has no codegen-participating union; the five
affected cases use that helper. Resolved.
<!-- SPARK_DEV_REVIEW_REPLY
{"feedback_id":"inline:4023166759","thread_id":"inline:4023166759","verdict_sha256":"6a099963095a21eb62dcca5d834d1640efeaefed2ee53a50abe468171a2023f3"}
-->
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]