Order a condition's conjuncts so a MobilitySpark call runs last - #56
Merged
estebanzimanyi merged 1 commit intoSep 16, 2026
Merged
estebanzimanyi merged 1 commit into
estebanzimanyi merged 1 commit into
Conversation
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
force-pushed
the
perf/order-join-conjuncts-by-cost
branch
from
September 16, 2026 09:25
b87b556 to
46ca891
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.
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.