[
https://issues.apache.org/jira/browse/SPARK-58902?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18109368#comment-18109368
]
AnhTris commented on SPARK-58902:
---------------------------------
I've been working on investigating this issue and have opened a PR with test
suite coverage and analysis: https://github.com/apache/spark/pull/58377
h3. Technical Investigation & Analysis: SPARK-58902
I have investigated this issue and confirmed the semantic inconsistency when
multi-referenced nondeterministic expressions reside within conditional
branches or join conditions. Below are the minimal reproduction steps, the root
cause analysis in Catalyst optimizer, and a detailed proposed fix direction
based on lazy per-row memoization.
----
h3. 1. Reproduction
The issue can be reproduced in Spark SQL whenever a nondeterministic common
expression (or expressions expanded by {{RuntimeReplaceable}} like {{BETWEEN}},
{{nullif}}, etc.) appears inside a conditional expression ({{CASE WHEN}} /
{{If}} / {{Coalesce}}):
{code:sql}
-- 1. Create a test table
CREATE OR REPLACE TEMPORARY VIEW t AS
SELECT id, CASE WHEN id % 2 = 0 THEN -1 ELSE 1 END AS a
FROM range(10);
-- 2. Query with BETWEEN on a nondeterministic generator inside CASE WHEN
SELECT id, a,
CASE WHEN a < 0 THEN false
ELSE randstr(3, 0) BETWEEN 'a' AND 'b' END AS in_range
FROM t;
{code}
*Observed Problem:*
* {{BETWEEN}} is syntactic sugar for {{(expr >= 'a' AND expr <= 'b')}} and
wraps the expression in {{With}}.
* Because {{With}} is inside the {{ELSE}} branch of {{CASE WHEN}},
{{RewriteWithExpression}} inlines {{randstr(3, 0)}} into each comparison
operand.
* As a consequence, {{randstr(3, 0)}} is evaluated twice per row in the
{{ELSE}} branch, yielding two distinct random strings for {{>= 'a'}} and {{<=
'b'}}.
* This violates the single-evaluation promise of {{With}} and produces
nondeterministic semantic corruption.
----
h3. 2. Root Cause
In Catalyst Optimizer ({{RewriteWithExpression.scala}}):
# {{RewriteWithExpression}} normally hoists multi-referenced {{With}}
definitions into a child {{Project}} to ensure exactly-once evaluation per row.
# Inside {{ConditionalExpression}} ({{CaseWhen}}, {{If}}), eager pre-evaluation
cannot be unconditionally placed into a child {{Project}} because an expression
that might raise an exception (such as division by zero, string cast error)
would execute eagerly on rows where the branch condition was false.
# Therefore, {{RewriteWithExpression}} falls back to inlining definitions:
{code:scala}
case c: ConditionalExpression =>
newExpr.transformUpWithPruning(_.containsPattern(WITH_EXPRESSION)) {
case With(child, defs) =>
val refToExpr = defs.map(d => d.id -> d.child).toMap
child.transformWithPruning(_.containsPattern(COMMON_EXPR_REF)) {
case ref: CommonExpressionRef => refToExpr(ref.id)
}
}
{code}
# Inlining duplicates evaluation of nondeterministic expressions ({{randstr}},
{{rand}}, {{uuid}}, {{uniform}}, {{shuffle}}, {{reflect}}).
# A similar issue occurs in join conditions when common expressions reference
columns from both sides of the join, forcing inlining with an explicit {{TODO}}
in {{RewriteWithExpression.scala}}.
----
h3. 3. Proposed Fix Direction
Introduce *lazy per-row memoization* for common expressions instead of
requiring eager pre-evaluation in a child {{Project}}:
# *Lazy Evaluable CommonExpressionRef:*
#* Make {{CommonExpressionRef}} evaluable at runtime by associating it with a
per-row cached state slot ({{evaluated: Boolean}}, {{value: Any}}, {{isNull:
Boolean}}).
#* On first reference during row evaluation, it evaluates its underlying
{{CommonExpressionDef}} child and stores the result in the slot.
#* Subsequent references in the same row read from the slot without
re-evaluating.
# *WholeStageCodegen & Interpreted Support:*
#* In codegen: Generate a row-scoped boolean flag {{eval_done}} and variable
cache. The computation is executed inside a guard {{if (!eval_done) { ...
eval_done = true; }}}.
#* In interpreted evaluation: Manage per-row evaluation cache in an execution
context.
# *Benefits:*
#* Safely preserves branch condition semantics without evaluating expressions
on unvisited rows.
#* Correctly guarantees single evaluation for nondeterministic expressions
across conditional branches, short-circuiting operators, and join conditions.
#* Eliminates fragile syntactic allowlists ({{canPreEvaluateInBranch}}).
> Evaluate a multi-referenced common expression lazily instead of
> pre-evaluating it in a Project
> ----------------------------------------------------------------------------------------------
>
> Key: SPARK-58902
> URL: https://issues.apache.org/jira/browse/SPARK-58902
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.0.3
> Reporter: Yang Jie
> Priority: Major
> Labels: pull-request-available
>
> `With` promises a common expression is evaluated only once even when
> referenced more
> than once (see the scaladoc on `With`). `RewriteWithExpression` keeps that
> promise by
> hoisting the definition into a `Project`, which works because a `Project` is
> evaluated
> for every row. Inside a conditional branch that is not true: the branch may
> not be
> evaluated at all, so the rule inlines instead, and inlining hands each
> reference its own
> evaluation.
> For a nondeterministic definition that is wrong, not just wasteful.
> SPARK-58818 fixed the
> case that is reachable through the generators which cannot raise, by
> pre-evaluating them
> in the child `Project` anyway. That fix is bounded by an allowlist
> (`canPreEvaluateInBranch`): pre-evaluation happens on every row, so an
> expression that can
> raise would turn a wrong result into a spurious error. Everything the
> allowlist turns down
> keeps the old inlining and the old wrong result, which the code says in as
> many words.
> h3. What is still wrong
> {code:sql}
> -- randstr is referenced twice by BETWEEN and inlined, so the two comparisons
> see
> -- two different strings
> SELECT CASE WHEN a < 0 THEN false
> ELSE randstr(3, 0) BETWEEN 'a' AND 'b' END
> FROM t
> {code}
> The same shape with {{nullif}}, and the same for {{uniform(lo, hi)}},
> {{shuffle(arr)}},
> {{reflect(...)}}, {{rand() / col}}, {{cast(rand() as int)}} -- anything
> nondeterministic
> that is not one of the six leaf generators on the allowlist. Note the
> allowlist
> approximates "cannot raise" by syntactic shape, so {{randstr(3, 0)}} is
> turned down for
> being a {{BinaryLike}} even though only {{randstr(-1, 0)}} raises.
> Outside a conditional branch these are correct, because the main rewrite
> hoists any
> multi-referenced non-cheap definition without consulting the allowlist. The
> wrong results
> are specific to a definition that is both inside a branch and turned down by
> the allowlist.
> There is a second, older instance of the same root cause. When a join
> condition's common
> expression references columns from both sides, no single child plan can hold
> the column,
> so `RewriteWithExpression` force-inlines it. The rule carries a TODO
> admitting this goes
> wrong for a nondeterministic definition and is kept only to match the old
> buggy behavior.
> h3. What lazy memoization would give
> An evaluable expression that caches its value per row, so a reference reads
> the value the
> first time it is reached and reuses it after. Then:
> * The allowlist can go. Evaluation happens only where the original expression
> would have
> evaluated, so an expression that can raise raises on exactly the rows it
> did before, and
> the wrong results above are fixed as one class rather than one generator at
> a time.
> * The evaluation domain is per reference reached, not per branch reached. A
> reference
> behind a short-circuiting operator (`a > 0 AND rand() BETWEEN 0.4 AND 0.6`)
> or inside a
> nested conditional is not read on every row of the branch; pre-evaluation
> cannot express
> that, memoization does.
> * The join TODO goes away, since the definition no longer has to be placed in
> a child
> plan.
> * The supporting machinery in `RewriteWithExpression` -- child projects,
> column naming,
> the per-child registry, projecting the extra columns away, the interaction
> with
> `CollapseProject` -- is no longer needed for these cases.
> h3. Why the existing machinery does not cover it
> Subexpression elimination cannot be reused.
> `EquivalentExpressions.updateExprInMap` is
> gated on `expr.deterministic`, so `rand`/`randn` never become common
> subexpressions, and
> `updateExprTree` additionally skips every `LeafExpression`, which excludes
> `uuid`,
> `monotonically_increasing_id`, `spark_partition_id` and `input_file_name`
> twice over.
> Where it does apply it is eager rather than lazy: codegen emits one `subExpr`
> function per
> common expression and calls them all up front, and
> `SubExprEliminationState.children` only
> orders dependencies. Interpreted evaluation does have per-row memoization in
> `SubExprEvaluationRuntime`, but its proxies come from the same gated map.
> So this needs a new expression with per-row caching plus codegen support.
> `With`, `CommonExpressionDef` and `CommonExpressionRef` are all `Unevaluable`
> today, which
> is why `RewriteWithExpression` has to eliminate them before execution.
> h3. A negative result worth recording
> SPARK-58818 also tried the closest thing reachable in the plan: keep the eager
> pre-evaluation but wrap the column in `If(branch is reached, definition,
> default)`, so a
> stateful generator is only advanced on the rows that reach the branch. It
> does not hold up.
> A guard in a plain `Project` is outside conditional evaluation, so
> subexpression
> elimination hoists a part repeated across guards to the top of the projection
> and evaluates
> it eagerly (`spark.sql.subexpressionElimination.skipForShortcutExpr` is false
> by default),
> and a condition that can raise then raises on rows the branch never reached:
> {code:sql}
> SELECT CASE WHEN a = 0 THEN false
> WHEN 6 / a > 2 THEN rand(1) BETWEEN 0 AND 1
> WHEN 6 / a < -2 THEN rand(2) BETWEEN 0 AND 1
> ELSE false END
> FROM t -- t contains a = 0
> {code}
> gives DIVIDE_BY_ZERO with the frame in `project_subExpr_0$`. Restricting
> guards to nodes
> that cannot raise fixes that but leaves only plain column comparisons: `a + b
> > 0`,
> `abs(a) > 0`, `upper(s) = 'X'`, `cast(a as long) > 0` all lose the guard,
> silently. The
> guard is a dead end; memoization is not affected by any of this, since it
> does not put the
> condition anywhere.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]