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

Reply via email to