[
https://issues.apache.org/jira/browse/SPARK-59603?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-59603:
-----------------------------------
Labels: pull-request-available (was: )
> Memoize a subexpression repeated inside a conditional branch
> ------------------------------------------------------------
>
> Key: SPARK-59603
> URL: https://issues.apache.org/jira/browse/SPARK-59603
> Project: Spark
> Issue Type: Improvement
> Components: SQL
> Affects Versions: 5.0.0
> Reporter: Yang Jie
> Priority: Major
> Labels: pull-request-available
>
> Subexpression elimination cannot reach a subexpression that is repeated
> inside a single branch of an {{if}} or {{case when}}, and the reason is
> structural rather than an oversight: 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 is in no group, so it is
> evaluated once per occurrence for every row that reaches the branch.
> Since SPARK-58818 a {{With}} is evaluable and memoizes per evaluation, which
> covers exactly this: the definition sits inside the branch, so nothing is
> computed for a row that does not reach it, and the references read one value.
> This adds an optimizer rule, {{MemoizeCommonExpressionsInBranches}}, that
> rewrites the tallest repeated subexpression of each conditionally evaluated
> child of an {{if}} / {{case when}} into a {{With}}. It runs late (after the
> simplification rules, so what it memoizes is what survives them) and leaves
> the {{With}} in the branch, which is the shape {{RewriteWithExpression}}
> keeps anyway.
> Measured on {{select case when id < 5 then udf(id) + udf(id) else 0 end from
> range(0, 10, 1, 1)}} with a counting UDF: 5 evaluations with the rule, 10
> without.
> Guarded by
> {{spark.sql.optimizer.memoizeCommonExpressionsInBranches.enabled}}, default
> false, so plans are unchanged unless it is turned on. A candidate is skipped
> when it is cheap ({{CollapseProject.isCheap}}), when it holds a
> common-expression reference or a lambda variable (neither can be evaluated
> where the definition would sit), when it holds an aggregate, window or
> generator expression (the planner looks for those where they stand), or when
> it holds a subquery expression. A branch that already holds a {{With}} is
> left alone, since {{RewriteWithExpression}} defers a nested {{With}} to a
> pass that no longer runs by then.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]