Compute each side of a nested-loop join once where it calls a MobilitySpark function - #57
Merged
estebanzimanyi merged 1 commit intoSep 16, 2026
Conversation
Member
Author
|
The diff carries the commit of branch perf/order-join-conjuncts-by-cost, which this branch builds on: that branch registers the Catalyst extension and the conjunct-ordering rule, and this one adds the second rule to the same extension. |
…ySpark function The Catalyst extension injects a second optimizer rule, MaterializeNestedLoopSides. Where a join carries no equality between its two sides, so Spark executes it as a nested loop over pairs of partitions, every side whose plan calls a user-defined function goes behind a repartition boundary. A join that has such an equality is left alone, since hashing reads each side once already; a side that carries no call is left alone, since recomputing a scan is what Spark is good at; and a side already behind a boundary is returned as it is, so the rule reaches a fixed point. The predicate that recognises a call is the one the conjunct-ordering rule uses. Without a boundary a nested loop recomputes a side for every partition of the other: with thirty-two partitions a side, each row is read, filtered and projected thirty-two times. Where that projection calls a user-defined function, as clipping a trajectory to a region does, the recomputation dominates the join, and the plan shows why, its two sides being chains of a scan, a filter and projections with nothing between them and the join. Measured, `mvn test` 6 of 6: the rule's test asserts a boundary under a nested-loop join whose side calls the function and none under a join by equality, its control runs the same query in a session without the extension and asserts no boundary appears, and the conjunct-ordering tests of the commit below continue to pass. Measured on a month of AIS positions, the proximity count over a week window of that month, one local session with a 16 GB heap and eight tasks: the query answers 278 in 37 s with both rules, against 776 s with the ordering rule alone and 1,210 s with neither. With adaptive execution disabled it answers in 35 s, so the boundary itself carries the difference rather than a replanning it makes possible. The plan carries the evidence beside the timing: with the rule an exchange stands under each side of the join, and without it each side is a chain of a scan, a filter and projections that the join reads again for every partition of the other side.
estebanzimanyi
force-pushed
the
perf/materialize-the-nested-loop-sides
branch
from
September 16, 2026 09:47
1c573b2 to
033ccc4
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Order a condition's conjuncts so a MobilitySpark call runs last
MobilitySpark registers a Catalyst extension,
org.mobilitydb.spark.catalyst.MobilitySparkExtensions,which a session enables by naming it in
spark.sql.extensions. It injects one optimizer rule,OrderConjunctsByCost: in a join condition and in a filter, every conjunct that calls a user-defined
function moves behind the conjuncts that do not, the relative order inside each group is kept, and a
condition already in that shape is returned unchanged.
ANDevaluates its operands in order andstops at the first false one, so a conjunct placed later is evaluated on fewer rows: the rewrite
only reduces evaluations and never introduces one, which preserves both the answer and the
exceptions a function can raise. A conjunct that is not deterministic keeps its place, since its
position is observable.
Spark attaches no cost to a user-defined function, so its optimizer leaves the order the query
produced. A proximity join over trajectories shows what that order costs: the plan evaluates the
distance between two trajectories on every pair the join enumerates, while the box and time
comparisons in the same conjunction run after it, and over a month of AIS positions those
comparisons admit 24,616 of the roughly 7.1 million pairs the join produces at the week window.
The extension is Java against Catalyst's Scala injection points, so the build gains no module: the
entry point extends scala.runtime.AbstractFunction1, the rule extends Rule, and the
plan walk uses a Java scala.runtime.AbstractPartialFunction.
Measured,
mvn test3 of 3: OrderConjunctsByCostTest asserts that a join condition naming thefunction first optimizes into one naming it last, that a condition already ordered is left alone,
and that each query answers what it answers without the rule; OrderConjunctsByCostControlTest runs
the same query in a session without the extension and asserts the call stays first, which is what
makes the first assertion a statement about the rule rather than about Spark's own optimizer.
Measured on a month of AIS positions, the proximity count over a week window of that month, one
local session with a 16 GB heap and eight tasks: the query answers 278 in 776 s with the rule and
1,210 s without it, over a floor of 497 s for the same plan counting the candidate pairs and
calling no distance at all. The distance falls from 713 s over roughly 7.1 million pairs to 279 s
over the 24,616 the comparisons admit, and the plan carries the evidence beside the timing, its
join condition ending with the call rather than leading with it.
Compute each side of a nested-loop join once where it calls a MobilitySpark function
The Catalyst extension injects a second optimizer rule, MaterializeNestedLoopSides. Where a join
carries no equality between its two sides, so Spark executes it as a nested loop over pairs of
partitions, every side whose plan calls a user-defined function goes behind a repartition
boundary. A join that has such an equality is left alone, since hashing reads each side once
already; a side that carries no call is left alone, since recomputing a scan is what Spark is good
at; and a side already behind a boundary is returned as it is, so the rule reaches a fixed point.
The predicate that recognises a call is the one the conjunct-ordering rule uses.
Without a boundary a nested loop recomputes a side for every partition of the other: with
thirty-two partitions a side, each row is read, filtered and projected thirty-two times. Where
that projection calls a user-defined function, as clipping a trajectory to a region does, the
recomputation dominates the join, and the plan shows why, its two sides being chains of a scan, a
filter and projections with nothing between them and the join.
Measured,
mvn test6 of 6: the rule's test asserts a boundary under a nested-loop join whoseside calls the function and none under a join by equality, its control runs the same query in a
session without the extension and asserts no boundary appears, and the conjunct-ordering tests of
the commit below continue to pass.
Measured on a month of AIS positions, the proximity count over a week window of that month, one
local session with a 16 GB heap and eight tasks: the query answers 278 in 37 s with both rules,
against 776 s with the ordering rule alone and 1,210 s with neither. With adaptive execution
disabled it answers in 35 s, so the boundary itself carries the difference rather than a
replanning it makes possible. The plan carries the evidence beside the timing: with the rule an
exchange stands under each side of the join, and without it each side is a chain of a scan, a
filter and projections that the join reads again for every partition of the other side.