Skip to content

Order a condition's conjuncts so a MobilitySpark call runs last - #56

Merged
estebanzimanyi merged 1 commit into
MobilityDB:mainfrom
estebanzimanyi:perf/order-join-conjuncts-by-cost
Sep 16, 2026
Merged

estebanzimanyi merged 1 commit into
MobilityDB:mainfrom
estebanzimanyi:perf/order-join-conjuncts-by-cost

Conversation

@estebanzimanyi

@estebanzimanyi estebanzimanyi commented Sep 16, 2026 •

Copy link
Copy Markdown
Member

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.

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<LogicalPlan>, 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.
@estebanzimanyi
estebanzimanyi force-pushed the perf/order-join-conjuncts-by-cost branch from b87b556 to 46ca891 Compare September 16, 2026 09:25
@estebanzimanyi
estebanzimanyi merged commit 9f5775d into MobilityDB:main Sep 16, 2026
2 checks passed
@estebanzimanyi
estebanzimanyi deleted the perf/order-join-conjuncts-by-cost branch September 16, 2026 09:45
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