LuciferYang opened a new pull request, #58882:
URL: https://github.com/apache/spark/pull/58882

   ### What changes were proposed in this pull request?
   
   Subexpression elimination cannot reach a subexpression that is repeated 
inside a single branch of an `if` or a `case when`, and the reason is 
structural: it evaluates its candidates before the projection, so 
`EquivalentExpressions` only collects what is always evaluated 
(`ConditionalExpression.alwaysEvaluatedInputs`) plus what every branch of a 
group shares (`branchGroups`, an intersection). A subexpression repeated inside 
one branch body and nowhere else belongs to no group, so each occurrence is 
evaluated for every row that reaches that branch.
   
   Since SPARK-58818 a `With` is evaluable and memoizes per evaluation, which 
is what this gap needs: the definition sits inside the branch, so nothing is 
computed for a row that takes another one, and the references read the one 
value.
   
   This adds `MemoizeCommonExpressionsInBranches`, which rewrites the tallest 
repeated subexpression of each conditionally evaluated child of an `if` / `case 
when` into a `With`. It runs after the simplification rules, so what it 
memoizes is what survives them, and it leaves the `With` in the branch -- the 
shape `RewriteWithExpression`, which runs right after `FinishAnalysis`, keeps 
anyway, so that rule does not have to run again.
   
   A candidate is skipped when reading it back would not be cheaper than 
recomputing it (`CollapseProject.isCheap`), when it holds a common-expression 
reference or a lambda variable -- neither can be evaluated where the definition 
would sit -- or when it holds a subquery expression. An aggregate, window or 
generator expression rules out the whole branch body rather than just a 
candidate holding one: the planner takes those out of the tree they stand in, 
and a reference left behind would be evaluated where its `With`, and so its 
definition, is no longer above it. A branch that already holds a `With` is left 
alone too, since `RewriteWithExpression` defers a nested `With` to a pass that 
no longer runs by this point.
   
   Being stateful is deliberately not a reason to refuse. A `ScalaUDF` is 
stateful because its encoder reuses an `UnsafeRow`, and an expensive UDF 
repeated in a branch is what this rule is for; the definition is evaluated once 
and read back within the same row, which is what a `With` left in a branch by 
`nullif(udf(x), 0)` already does. What would be unsafe is a value that changes 
per evaluation, and those are nondeterministic -- `EquivalentExpressions` never 
records one.
   
   ### Why are the changes needed?
   
   Over `select case when id < 5 then udf(id) + udf(id) else 0 end from 
range(0, 10, 1, 1)` with a counting UDF, the definition is evaluated 5 times 
with the rule and 10 times without: once for each of the five rows that reach 
the branch, rather than once per reference. Neither subexpression elimination 
nor `RewriteWithExpression` covers that today, for the reason above.
   
   The rule is behind 
`spark.sql.optimizer.memoizeCommonExpressionsInBranches.enabled`, default 
false, so it changes no plan unless it is turned on. Two things are 
deliberately left for later and are the reason for the default: the planning 
cost has not been measured (one `EquivalentExpressions` per branch body), and 
the golden plans of TPCDS and `PlanStability*` would move once a `With` reaches 
physical plans at scale, which deserves its own evaluation.
   
   ### Does this PR introduce _any_ user-facing change?
   
   No. The config is off by default; with it on, results are unchanged and 
`EXPLAIN` shows a `With` inside the branch.
   
   ### How was this patch tested?
   
   New `MemoizeCommonExpressionsInBranchesSuite`: the repeated subexpression 
becomes one definition read by two references and the `With` stays inside the 
branch; the plan is untouched while the config is off; a cheap subexpression, 
an always-evaluated first condition, a body holding an aggregate expression, 
and a branch that already holds a `With` are all left alone.
   
   A new case in `ColumnExpressionSuite` counts the evaluations through an 
accumulator UDF -- 5 with the rule, 10 without -- and asserts the rows and 
whether a `With` survived into the optimized plan, so that "the rule did not 
fire" and "the memoization did not take" fail differently. It uses one 
`collect()` rather than `checkAnswer`, which runs the plan more than once and 
would sum the counts.
   
   catalyst `optimizer.*` (1492 passed, 1 ignored) and `ColumnExpressionSuite` 
+ `SubexpressionEliminationSuite` (162 passed) with the default. With the 
default temporarily flipped to true, `DataFrameFunctionsSuite` + 
`ColumnExpressionSuite` + `DataFrameAggregateSuite` (499 passed) and catalyst 
`optimizer.*` (1492 passed) also pass. `catalyst/scalastyle`, 
`catalyst/Test/scalastyle` and `sql/Test/scalastyle` clean. No benchmark was 
run.
   
   ### Was this patch authored or co-authored using generative AI tooling?
   
   Generated-by: Claude Opus 5
   


-- 
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]

Reply via email to