Skip to content

Compute each side of a nested-loop join once where it calls a MobilitySpark function - #57

Merged
estebanzimanyi merged 1 commit into
MobilityDB:mainfrom
estebanzimanyi:perf/materialize-the-nested-loop-sides
Sep 16, 2026
Merged

estebanzimanyi merged 1 commit into
MobilityDB:mainfrom
estebanzimanyi:perf/materialize-the-nested-loop-sides

Conversation

@estebanzimanyi

Copy link
Copy Markdown
Member

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. AND evaluates its operands in order and
stops 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 test 3 of 3: OrderConjunctsByCostTest asserts that a join condition naming the
function 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 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

Copy link
Copy Markdown
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
estebanzimanyi force-pushed the perf/materialize-the-nested-loop-sides branch from 1c573b2 to 033ccc4 Compare September 16, 2026 09:47
@estebanzimanyi
estebanzimanyi merged commit 432ebc3 into MobilityDB:main Sep 16, 2026
2 checks passed
@estebanzimanyi
estebanzimanyi deleted the perf/materialize-the-nested-loop-sides branch September 16, 2026 09:52
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant