[
https://issues.apache.org/jira/browse/SPARK-60040?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Luka Zdravic updated SPARK-60040:
---------------------------------
Description:
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.
Assigned to : [~lukazdravic]
was:
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.
> 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
> Priority: Minor
>
> 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.
> Assigned to : [~lukazdravic]
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]