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]
