Luka Zdravic created SPARK-60040:
------------------------------------
Summary: Carry the AS-OF join direction to SortMergeAsOfJoinExec
instead of inferring it from the as-of condition
Key: SPARK-60040
URL: https://issues.apache.org/jira/browse/SPARK-60040
Project: Spark
Issue Type: Improvement
Components: SQL
Affects Versions: 4.4.0
Reporter: Luka Zdravic
SortMergeAsOfJoinExec infers the AS-OF direction from the shape of
asOfCondition:
{code:scala}
private val asOfDirection: AsOfJoinDirection =
splitConjunctivePredicates(asOfCondition).head match {
case c: BinaryComparison
if c.left.references.subsetOf(left.outputSet) &&
c.right.references.subsetOf(right.outputSet) =>
c match {
case _: GreaterThanOrEqual | _: GreaterThan => Backward
case _: LessThanOrEqual | _: LessThan => Forward
case _ => Nearest
}
case _ => Nearest
}
{code}
The logical plan knows the direction exactly, but does not keep it:
* DataFrame path: AsOfJoin.apply takes a direction, but uses it only to build
asOfCondition and orderExpression.
* SQL path: ResolveAsOfJoin and AsOfJoinResolver set matchOperator = None after
resolution.
So the exec node relies on the optimized condition keeping the "left-side op
right-side" shape. If a rewrite changes that shape, the direction silently
becomes Nearest. Example: a backward condition l.ts >= r.ts rewritten as r.ts
<= l.ts falls through to case _ => Nearest. A Nearest scan can return a
different row than a Backward or Forward scan, and it computes a distance that
Backward and Forward do not need.
h3. Proposed fix
* Add direction: AsOfJoinDirection to AsOfJoin. Set it in AsOfJoin.apply and in
the SQL resolution paths.
* Pass it to SortMergeAsOfJoinExec and match on it exhaustively. Remove the
shape-based inference.
* Add a test where the as-of condition does not have the "left op right" shape
(for example, right side first), and check that the join keeps the requested
direction.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]