andygrove opened a new issue, #5824:
URL: https://github.com/apache/datafusion-comet/issues/5824
### 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)`:
```sql
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:
```sql
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.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]