cloud-fan commented on code in PR #58419:
URL: https://github.com/apache/spark/pull/58419#discussion_r3999365352
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/basicPhysicalOperators.scala:
##########
@@ -976,11 +977,13 @@ case class UnionExec(children: Seq[SparkPlan]) extends
SparkPlan with CodegenSup
}
}
- override def outputPartitioning: Partitioning = {
- if (!conf.getConf(SQLConf.UNION_OUTPUT_PARTITIONING)) {
- return super.outputPartitioning
- }
-
+ /**
+ * The SPARK-52921 pass-through partitioning, derived from the children.
`isPlainUnion` answers
Review Comment:
**Nit (P3):** These are not equivalent: when `UNION_OUTPUT_PARTITIONING` is
disabled, `rawPartitioning` may still return a concrete child-derived
partitioning while `isPlainUnion` is true and `outputPartitioning` reports
unknown. Please describe `rawPartitioning` as the child-derived candidate and
call out the independent configuration-disabled branch.
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/basicPhysicalOperators.scala:
##########
@@ -1035,13 +1038,78 @@ 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, 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) ||
+ rawPartitioning.isInstanceOf[UnknownPartitioning]
+ }
+
+ private def stampedDecisions: Option[UnionExec.Decisions] =
+ getTagValue(UnionExec.DECISIONS)
+
+ /**
+ * 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. 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(): Unit = if
(stampedDecisions.isEmpty) {
+ setTagValue(UnionExec.DECISIONS, UnionExec.Decisions(
+ plainUnion = isPlainUnion,
+ unionCodegenEnabled =
conf.getConf(SQLConf.WHOLESTAGE_UNION_CODEGEN_ENABLED),
+ maxChildren = conf.getConf(SQLConf.WHOLESTAGE_UNION_MAX_CHILDREN)))
+ }
+
+ /**
+ * 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, where re-deriving at the gate would
have fused it.
Review Comment:
**Nit (P3):** `where` is contrastive here; use `whereas`.
##########
sql/core/src/test/scala/org/apache/spark/sql/execution/UnionCodegenSuite.scala:
##########
@@ -628,6 +659,290 @@ 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.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 where the plain-union
decision is stamped, 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;
`unionInsideWSCG` would also
+ // match a union that an `InputAdapter` left inside the stage unfused,
which is exactly
+ // the degradation this guard has to catch.
+ 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. This assertion is what fails there, because nothing
forces the copy's reason
+ // before it. It has to sit after the flip, as it does here: taken
while the conf was still
+ // on, it would warm a memoizing implementation with the answer this
test needs it not to
+ // have.
+ val copy = fusedUnions(planned)
+ assert(copy.size == 1)
+ 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.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")
+ 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 prepared union keeps its layout when nothing read it
during preparation") {
+ // With whole-stage codegen off, no gate consults the union while the plan
is prepared, and a
+ // root union has no parent to ask for its partitioning either. First-read
initialization would
+ // then decide at execution, under whatever the conf says by then;
`StampUnionDecisions` decides
+ // during preparation instead.
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.WHOLESTAGE_CODEGEN_ENABLED.key -> "false") {
+ val plan = withSQLConf(SQLConf.UNION_OUTPUT_PARTITIONING.key -> "true") {
+ 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")))
+ .queryExecution.executedPlan
+ }
+ withSQLConf(SQLConf.UNION_OUTPUT_PARTITIONING.key -> "false") {
+ // Co-partitioned children pass their four partitions through; a plain
concatenation would
+ // report eight.
+ assert(plan.execute().getNumPartitions == 4,
+ "a prepared union must execute by the layout it was prepared with")
+ }
+ }
+ }
+
+ test("SPARK-59122: a later stamping pass fills in a fresh union and keeps
stamped ones") {
Review Comment:
**Non-blocking (P2):** This calls `StampUnionDecisions` directly, so it
remains green if any of the new post-extension barriers is removed or
misordered. Please add integration coverage that creates a `UnionExec` through
real `SparkSessionExtensions` for the classic columnar, AQE query-stage
preparation, and AQE post-stage columnar paths; those are the wiring points
whose regression would recreate the decision/metrics mismatch.
**Recommended change:** Add extension-driven regression coverage that
creates a fresh UnionExec from the non-AQE columnar hook and from the AQE
query-stage preparation and post-stage optimization/columnar hooks, then
observes a configuration change after preparation to prove each node retained
the decision taken by its post-hook barrier.
**Why this works:** Use the existing SparkSessionExtensions test-session
pattern and injected Rule or ColumnarRule implementations to introduce a
semantically valid UnionExec at each late phase. Make every case discriminate
the barrier by changing a captured union configuration afterward and asserting
the extension-created node's final layout or codegen/metric behavior; the case
must fail if its corresponding StampUnionDecisions entry is deleted or moved
before the hook.
**Scope:** sql/core/src/test/scala/org/apache/spark/sql
**Compatibility:** The preparation-time snapshot, tag propagation for
intended copies, child-derived codegen gates, and all documented
result-preserving optimization tradeoffs remain unchanged.
**Risks:** A test can pass vacuously if an earlier barrier stamps the
original plan rather than the UnionExec created by the targeted hook. An
injected rule that duplicates or reshapes rows without a fixed oracle can turn
the test into a plan-shape assertion rather than a lifecycle regression test.
**Constraints:** Use supported SparkSessionExtensions hooks and a real
preparation or AQE stage lifecycle. Ensure each extension creates a fresh
tagless UnionExec after the barrier whose placement is under test. Keep
production rule ordering and write-once decision semantics unchanged.
**Success:** Removing or moving the post-columnar StampUnionDecisions entry
in QueryExecution makes a focused non-AQE case fail. Removing or moving the
barrier after AQE queryStagePrepRules makes a focused AQE case fail. Removing
or moving the barrier after AQE stage-optimizer or columnar processing makes a
focused AQE stage case fail. The existing cache, configuration-flip, and
write-once regression cases continue to pass.
--
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]