[ 
https://issues.apache.org/jira/browse/SPARK-59614?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-59614:
-----------------------------------
    Labels: pull-request-available  (was: )

> Compile the AS-OF join as-of/residual conditions to a Predicate
> ---------------------------------------------------------------
>
>                 Key: SPARK-59614
>                 URL: https://issues.apache.org/jira/browse/SPARK-59614
>             Project: Spark
>          Issue Type: Improvement
>          Components: SQL
>    Affects Versions: 4.1.0
>            Reporter: David Mollitor
>            Priority: Minor
>              Labels: pull-request-available
>
> h2. Problem
> {{{}SortMergeAsOfJoinScanner}({{sql/core/.../execution/joins/SortMergeAsOfJoinExec.scala{}}})
>  has no whole-stage codegen, so its inner right-buffer scan is always 
> interpreted. For every buffered right row it evaluated the as-of condition 
> with {{boundAsOfCond.eval(joinedRow)}} on a bound {{{}Expression{}}}, and the 
> residual condition with {{{}boundResidualCond.forall(...){}}}. Interpreted 
> {{Expression.eval}} walks the expression tree, and {{BoundReference.eval}} 
> boxes every operand it reads into non-cached {{java.lang.Double}} / 
> {{java.lang.Long}} objects, which allocate.
> JFR profiling of {{AsOfJoinBenchmark}} (sort-merge cases isolated) showed 
> this dominated the scan:
>  * CPU ~24%: {{boxToInteger}} 17.6% + {{Double.valueOf}} 4.1% + 
> {{Long.valueOf}} 2.0%
>  * Allocation ~33%: {{java.lang.Double}} 20% + {{java.lang.Long}} 13%
>  * Stacks trace to {{boxToInteger -> BoundReference.eval -> 
> BinaryExpression.eval -> findBestBackwardForward}}
> h2. Change
> Compile the as-of and residual conditions to {{BasePredicate}} via 
> {{Predicate.create(...)}} and evaluate them through the primitive-returning 
> {{eval(InternalRow): Boolean}} instead of a boxed interpreted 
> {{{}Expression.eval{}}}:
>  * {{boundAsOfCond}} becomes a {{{}BasePredicate{}}}; {{boundResidualCond}} 
> becomes {{{}Option[BasePredicate]{}}}.
>  * * Each predicate is initialized per partition 
> ({{{}initialize(partitionIndex){}}}, index from
> {{{}TaskContext.getPartitionId(){}}}), matching the {{CartesianProductExec}} 
> idiom.
>  * Call sites use {{if (boundAsOfCond.eval(joinedRow))}} plus a small 
> {{residualHolds}} helper that avoids the {{Option.forall}} closure and its 
> boxed result.
>  * The {{orderExpression}} (a distance value, not a boolean) intentionally 
> stays a bound {{{}Expression{}}}.
> This removes both the interpreted tree-walk and the per-operand boxing on the 
> hot scan.



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