dwsmith1983 opened a new issue, #6787:
URL: https://github.com/apache/datafusion-comet/issues/6787

   ### Describe the bug
   
   A native hash join can return its streamed side's rows out of order while 
Spark relies on that order, so a sort Spark removed as redundant is lost and 
`sortWithinPartitions` returns unsorted partitions.
   
   Spark's `HashJoin.outputOrdering` is the streamed side's ordering, and 
`CometHashJoinExec` and `CometBroadcastHashJoinExec` copy it from the Spark 
join (`op.outputOrdering` in `operators.scala`). When the streamed side already 
arrives partitioned and sorted on the key, with no exchange in between, 
`RemoveRedundantSorts` drops a sort above the join on that key.
   
   Natively the join builds on the left, so a `BuildRight` join becomes a 
DataFusion join whose probe side is the streamed side. DataFusion keeps the 
probe side's order for unmatched rows only when that input declares an 
ordering: `HashJoinExec` passes `self.right.output_ordering().is_some()` 
(`datafusion-physical-plan` 55.1.0, `joins/hash_join/exec.rs:1574`) to 
`append_right_indices`, which otherwise appends the unmatched probe rows after 
the matched ones in each batch. A sorted input that reaches the native plan 
from the JVM goes through `ScanExec`, which declares no ordering, so for a left 
outer join building on the right the unmatched rows come out after the matched 
ones.
   
   When the sort runs natively in the same plan as the join, DataFusion sees 
its ordering and the output stays sorted.
   
   ### Steps to reproduce
   
   Parquet tables `big` as `(i % 100, i)` for 10000 rows and `small` as `(i * 
10, i)` for 10 rows, with `spark.comet.convert.inMemoryCache.enabled=true`, 
`spark.sql.adaptive.enabled=false`, `spark.sql.autoBroadcastJoinThreshold=-1` 
and `spark.sql.shuffle.partitions=2`:
   
   ```scala
   val streamed = spark.table("big").repartition(2, 
$"_1").sortWithinPartitions("_1").cache()
   streamed.count()
   val build = spark.table("small")
   val df = streamed
     .join(build.hint("shuffle_hash"), streamed("_1") === build("_1"), 
"left_outer")
     .select(streamed("_1").as("k"), streamed("_2").as("v"), 
build("_2").as("w"))
     .sortWithinPartitions("k")
   df.queryExecution.toRdd
     .mapPartitions { it =>
       val keys = it.map(_.getInt(0)).toArray
       Iterator(keys.sameElements(keys.sorted))
     }
     .collect()
   // Array(false, false); the rows match Spark's
   ```
   
   The plan has no sort above the join:
   
   ```
   CometProject [k, v, w]
   +- CometHashJoin [_1], [_1], LeftOuter, BuildRight
      :- CometSparkColumnarToColumnar
      :  +- InMemoryTableScan [_1, _2]
      +- CometExchange hashpartitioning(_1, 2)
   ```
   
   `build.hint("broadcast")` gives the same result with 
`CometBroadcastHashJoin`. Reproduced on `main` (9876862df) with Spark 4.1. 
`branch-1.1` copies the Spark join's output ordering the same way.
   
   ### Expected behavior
   
   Every partition sorted by `k`, as in Spark.
   
   ### Additional context
   
   It needs a sorted streamed side that reaches the native join from outside 
the native plan. Here that is a cached table read through 
`spark.comet.convert.inMemoryCache.enabled`; I have not checked which other 
inputs reach it, such as Comet's native cache scan or another JVM operator 
below a native join. A fix could have the two hash join execs report no output 
ordering for a join type and build side whose unmatched probe rows DataFusion 
may reorder, unless the probe child is known to keep its order in the native 
plan, or have the scan declare the ordering it receives.
   
   Found while working on #6770.
   


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