Skip to content

Comet join and explode operators omit joinType / outer from equals, so exchange reuse returns wrong results #5824

Description

@andygrove

Describe the bug

Several Comet physical operators hand-write equals/hashCode so that nativeOp, originalPlan
and serializedPlanOpt stay out of plan identity. Two of those overrides leave out a field that
changes the operator's results, so plans that compute different things canonicalize as equal and
ReuseExchangeAndSubquery shares a shuffle between them. The query then returns one branch's rows
twice.

This is the same defect that #5470 just fixed for CometHashAggregateExec (which omitted
resultExpressions). Two more instances are still live:

1. joinType is missing from all three join operators.
CometHashJoinExec (operators.scala:2450), CometBroadcastHashJoinExec (2597) and
CometSortMergeJoinExec (2789) each declare joinType: JoinType as a constructor field, list it
in stringArgs, and omit it from both equals and hashCode.
CometBroadcastNestedLoopJoinExec (2237) does include it, which suggests the other three were
simply missed rather than deliberately excluded.

For most join-type pairs the output comparison rescues equality, because nullability differs.
LeftSemi and LeftAnti are the exception: identical output, identical keys, identical condition,
identical build side. They canonicalize to the same plan.

Normally InferFiltersFromConstraints adds isnotnull(key) to the semi join's left child and not
to the anti join's, so the subtrees differ and the collision stays hidden. Writing the null check
explicitly in both branches removes that incidental protection.

2. CometExplodeExec never captures GenerateExec.outer.
convert writes op.outer into the protobuf via .setOuter(op.outer) (operators.scala:1490), but
createExec (1498) does not carry it onto the case class, so equals has no way to see it and
nativeOp is excluded by design. output does not disambiguate either: explode_outer forces the
generator output nullable, and plain explode over an array<int> with containsNull = true is
already nullable.

InferFiltersFromGenerate masks this for the common case by adding size(arr) > 0 AND isnotnull(arr) under non-outer generators. But that rule bails out when the generator input is not
a bare Attribute (Optimizer.scala:1705), so explode(s.arr), explode(split(...)),
explode(slice(...)) and friends get no filter and the collision is reachable on stock config.

Steps to reproduce

Both reproduce on Spark 4.1 / Scala 2.13 / JDK 17 with default configuration, against
c8ee6aef5. Comet native shuffle enabled, exchange reuse at its default of on.

Joins. Two Parquet tables, l = (0,10), (1,11), (2,12) and r = (0,100), (1,101):

SELECT _1, _2 FROM l WHERE _1 IS NOT NULL AND EXISTS     (SELECT 1 FROM r WHERE r._1 = l._1)
-- .repartition(2, col("_2")) on each branch, then UNION ALL with:
SELECT _1, _2 FROM l WHERE _1 IS NOT NULL AND NOT EXISTS (SELECT 1 FROM r WHERE r._1 = l._1)

Repartitioning on _2 rather than the join key matters, otherwise EnsureRequirements optimizes
the shuffle out and there is nothing to reuse.

join hint Spark Comet
SHUFFLE_HASH [0,10] [1,11] [2,12] [0,10] [0,10] [1,11] [1,11]
MERGE [0,10] [1,11] [2,12] [0,10] [0,10] [1,11] [1,11]
BROADCAST [0,10] [1,11] [2,12] [0,10] [0,10] [1,11] [1,11]

The Comet plan replaces the whole NOT EXISTS branch with a reuse of the EXISTS shuffle:

CometUnion Union, [_1#4, _2#5]
:- CometExchange hashpartitioning(_2#5, 2), REPARTITION_BY_NUM, CometNativeShuffle, [plan_id=362]
:  +- CometHashJoin [_1#4], [_1#6], LeftSemi, BuildRight
:     :- CometExchange hashpartitioning(_1#4, 2), ENSURE_REQUIREMENTS, CometNativeShuffle
:     :  +- CometFilter [_1#4, _2#5], isnotnull(_1#4)
:     :     +- CometNativeScan parquet [_1#4,_2#5]
:     +- CometExchange hashpartitioning(_1#6, 2), ENSURE_REQUIREMENTS, CometNativeShuffle
:        +- CometFilter [_1#6], isnotnull(_1#6)
:           +- CometNativeScan parquet [_1#6]
+- ReusedExchange [_1#8, _2#9], CometExchange hashpartitioning(_2#5, 2), REPARTITION_BY_NUM, [plan_id=362]

Vanilla Spark reuses only the two scan-side exchanges and correctly declines the top one.
CometHashJoinExec.sameResult returns true across the semi/anti pair, and every term in equals
(output, leftKeys, rightKeys, condition, buildSide, left, right, serializedPlanOpt)
compares equal field by field.

Explode. A Parquet table t(k int, arr array<int>, s struct<arr: array<int>>) with rows
(1, [10,20], ...), (2, [], ...), (3, null, ...), and each branch repartitioned on k before
the union:

SELECT k, explode(s.arr)       AS v FROM t
-- UNION ALL
SELECT k, explode_outer(s.arr) AS v FROM t
generator input reused exchanges Comet / Spark Spark Comet
arr (bare attribute) 0 / 0 6 rows 6 rows, masked by the inferred filter
arr, InferFiltersFromGenerate excluded 1 / 0 6 rows 4 rows
s.arr (struct field), stock config 1 / 0 6 rows 4 rows
slice(arr, 1, 10), stock config 1 / 0 6 rows 4 rows

Comet drops (2, null) and (3, null); the entire explode_outer branch becomes a
ReusedExchange pointing at the explode branch.

Expected behavior

Comet should return the same rows as Spark. Two plans that compute different results should not
compare equal, so exchange reuse should not fire across them.

Additional context

The fix looks small in both cases: add joinType to equals and hashCode on the three join
operators, and add an outer: Boolean field to CometExplodeExec and include it in equals and
hashCode. Each wants a regression modelled on the ones added in #5470.

Beyond the two instances, it would be worth adding a guard so the next operator does not repeat
this. A test that reflects over every CometNativeExec subclass and asserts that each constructor
parameter is either referenced by equals or named on an explicit exclusion list (nativeOp,
originalPlan) would have caught all three of these, including the aggregate case, before they
shipped.

Worth noting for whoever picks this up: in both cases an unrelated optimizer rule normally makes
the two subtrees differ, which is why this has gone unnoticed. A regression test has to defeat that
masking deliberately, either with an explicit IS NOT NULL on both join branches or with a
non-Attribute generator input.

Found while doing a post-merge review of #5470.

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

    area:joinsJoin operators and dynamic filter pushdownbugSomething isn't workingcorrectnesspriority:criticalData corruption, silent wrong results, security issuesrequires-triage

    Type

    No type

    Projects

    No projects

      Milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions