Skip to content

KNN joins silently drop additional ON conditions in broadcast and multi-conjunct plans #3398

Description

@jiayuasu

Additional non-equality join conditions are ignored by a query-side broadcast KNN join. A regular KNN join also ignores them when the optimized condition is (ST_KNN AND predicate1) AND predicate2. The same single predicate is correctly applied by the regular KNN plan, so changing the physical strategy or adding another conjunct changes query results.

Reproduced with Spark 4.1.2, Java 17, and Sedona 2.0.0-SNAPSHOT. The tested KNN planner/execution sources match master da6b11c9d2.

Reproduction

With an initialized Sedona Spark session, create these relations from Python rows to prevent constant folding of the payload comparisons:

spark.conf.set("spark.sql.adaptive.enabled", "false")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
spark.conf.set("spark.sedona.join.autoBroadcastJoinThreshold", "-1")
spark.conf.set("spark.sedona.join.knn.includeTieBreakers", "false")

schema = "id int, x double, y double, score int, cat int"
a = spark.createDataFrame([
    (1, 0.0, 0.0, 10, 100),
    (2, 10.0, 10.0, 10, 100),
], schema)
b = spark.createDataFrame([
    (101, 1.0, 1.0, 20, 100),
    (102, 11.0, 11.0, 20, 100),
], schema)
for name, frame in [("a", a), ("b", b)]:
    frame.selectExpr("id", "ST_Point(x,y) AS g", "score", "cat").createOrReplaceTempView(name)

Every query-side score is 10 and every candidate-side score is 20. Therefore every query containing a.score > b.score must return no rows, regardless of how KNN and the extra condition are ordered semantically.

Run each query with spark.sql(sql).collect() and inspect spark.sql(sql).explain().

-- Control: returns [(1, 101), (2, 102)].
SELECT a.id AS qid, b.id AS oid
FROM a JOIN b ON ST_KNN(a.g, b.g, 1, false);

-- Control: correctly returns [].
SELECT a.id AS qid, b.id AS oid
FROM a JOIN b ON ST_KNN(a.g, b.g, 1, false)
  AND a.score > b.score;

-- Bug 1: expected []; actual [(1, 101), (2, 102)].
SELECT /*+ BROADCAST(a) */ a.id AS qid, b.id AS oid
FROM a JOIN b ON ST_KNN(a.g, b.g, 1, false)
  AND a.score > b.score;

-- Bug 2: expected []; actual [(1, 101), (2, 102)].
SELECT a.id AS qid, b.id AS oid
FROM a JOIN b ON ST_KNN(a.g, b.g, 1, false)
  AND a.score > b.score AND a.id < b.id;

The observed physical operators are below, with expression IDs omitted and Sedona's verbose expression class name shortened to ST_KNN:

Single-conjunct control:
KNNJoin ..., Inner, 1, false,
  (ST_KNN AND (a.score > b.score)), (a.score > b.score)

Query-side broadcast:
BroadcastQuerySideKNNJoin ..., LeftSide, Inner, 1, false, KNN, false

Regular two-conjunct query:
KNNJoin ..., Inner, 1, false,
  ((ST_KNN AND (a.score > b.score)) AND (a.id < b.id))

The optimized logical plans still contain the extra predicates. The faulty physical plans have no residual filter above the join; the regular plan retains the original condition for display but lacks the separate executable extra condition present in the working control.

Likely cause

In JoinQueryDetector.scala:

  • planBroadcastJoin constructs both broadcast KNN executors with condition = null and extraCondition = None (lines 982–983 and 997–998).
  • planKNNJoin calls extractExtraKNNJoinCondition (line 844). That helper checks only whether either immediate child of the outer And is ST_KNN (lines 906–921). For the observed left-associated two-conjunct condition, neither immediate child is ST_KNN, so it returns None.

Both paths should retain and execute all additional conjuncts without evaluating the ST_KNN marker as a scalar expression. Regression tests should cover one and multiple non-equality conditions in regular and both broadcast KNN plans, including conjunct ordering.

Reproduction limits

Using a.cat = b.cat instead is not a reliable reproduction of silent predicate loss under the default spatial join optimization mode. In this environment it selects SortMergeJoin without a hint or BroadcastHashJoin with BROADCAST(b); both executions fail with KNN predicate is not supported when evaluating the remaining ST_KNN condition. The inequalities above reliably select KNN plans.

The candidate-side broadcast variant (BROADCAST(b) with a.score > b.score) selected BroadcastObjectSideKNNJoin with the condition absent, confirming the planner omission, but its row output was not verified: this local runtime encountered an unrelated JTS IndexSerde/AbstractSTRtree classloader IllegalAccessError. The two silent wrong-result cases above were both executed successfully.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions