dwsmith1983 opened a new pull request, #6810:
URL: https://github.com/apache/datafusion-comet/pull/6810

   ## Which issue does this PR close?
   
   Closes #6787.
   
   ## Rationale for this change
   
   Spark's `HashJoin.outputOrdering` is the streamed side's ordering, and the 
Comet hash join execs report it too. When the streamed side already arrives 
sorted, Spark removes a sort above the join on that basis. The native join does 
not always keep that order:
   
   - Comet runs `LeftOuter` with `BuildRight` and `RightOuter` with `BuildLeft` 
as a DataFusion `Right` join, which keeps unmatched probe rows in order only 
when it sees the probe input sorted (`datafusion-physical-plan` 55.1.0, 
`hash_join/exec.rs:1574`). When the order comes from outside the native plan, 
such as a cached sorted table, the unmatched rows come after the matched ones 
in each batch.
   - A null-aware anti join (`NOT IN`) runs unswapped, so DataFusion builds its 
hash table on Spark's streamed side and concatenates those batches in reverse 
(`exec.rs:2344`, `2365`). Across batches the output comes out in reverse order.
   
   `sortWithinPartitions` and `SORT BY` above such a join then return unsorted 
partitions.
   
   ## What changes are included in this PR?
   
   - `HashJoin` gets an `output_ordering` field. The serde fills it with 
Spark's reported ordering, bound to the join output, for the three shapes 
above, and only when that ordering is non-empty. It falls back to Spark when 
the ordering has a type the native sort does not support or an expression that 
cannot be serialized.
   - The native planner adds a `SortExec` above the hash join when the join's 
equivalence properties do not already satisfy that ordering. This follows the 
sort aggregate handling of `ordered_by_grouping_keys`. When the streamed side 
is sorted inside the same native plan, for example by a native sort or a 
sort-merge join, DataFusion keeps the order and no sort is added. Float keys 
are normalized the way a native sort below would normalize them, so that 
comparison holds.
   - The compatibility guide gets a hash join note, and the operator table 
lists the new fallback.
   
   The sort runs only when the order comes from the JVM, which takes a 
non-default conversion such as `spark.comet.convert.inMemoryCache.enabled` or 
`spark.comet.sparkToColumnar.enabled`, or for a `NOT IN` join whose streamed 
side reports an ordering. There it replaces unsorted output. The sort registers 
two memory consumers in the task, so under the fair pool the build side's share 
is smaller while it runs.
   
   `branch-1.1` has the same code path, so this is a backport candidate.
   
   **Merge order:** #6785 should merge first. Both change `CometJoinSuite` near 
the same tests, and once #6785 lands I'll merge main here and resolve that. 
Draft #6437 also takes `HashJoin` field 11; whichever merges second renumbers.
   
   ## How are these changes tested?
   
   `CometJoinSuite`, checking rows against Spark and every partition's order 
with Spark's own ordering:
   
   - `shuffle_hash` and `broadcast`, `left_outer` (BuildRight) and 
`right_outer` (BuildLeft), over a cached sorted streamed side, with Spark's 
sort above the join removed; plus a join condition, AQE on, and a small batch 
size so the sort merges several batches. The plan has no sort above the join, 
the join class matches the hint, and the join's build metrics still reach the 
Spark node.
   - A descending, nulls last order with null keys; a double with NaN and -0.0; 
a two-key order.
   - A streamed side sorted by a native sort below the join.
   - A `NOT IN` join over a cached sorted side whose build spans several 
batches.
   - Only the affected shapes send an ordering (an inner join sends none), and 
an unsupported sort type or expression falls back with its reason.
   
   Planner tests: the two outer shapes get a sort and return rows in order; a 
probe sorted natively, by a sort or a sort-merge join, gets none, including on 
a float key; no ordering adds no sort.
   
   The cached-side and `NOT IN` tests fail on main. Removing the shape 
restriction, the field, either fallback guard, the `ordering_satisfy` check or 
the float normalization each fails a test. `CometJoinSuite` passes on Spark 
3.4, 3.5, 4.0 and 4.1, Spark 4.2 compiles, and the strict warnings build 
passes. The Spark SQL `sql_core-1` shard (`dev/local-ci.sh`, Spark 4.1.3), 
which runs `RemoveRedundantSortsSuite` and the inner, outer, existence and hint 
join suites, passes.
   


-- 
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]

Reply via email to