SubhamSinghal commented on code in PR #25840:
URL: https://github.com/apache/datafusion/pull/25840#discussion_r4178690446


##########
datafusion/physical-plan/src/joins/piecewise_merge_join/classic_join.rs:
##########
@@ -631,18 +588,121 @@ fn build_matched_indices_and_mark_buffered(
     )?)
 }
 
-// Creates a record batch from the unmatched indices on the streamed side
-fn create_unmatched_batch(
-    streamed_indices: &mut PrimitiveBuilder<UInt32Type>,
-    stream_batch: &SortedStreamBatch,
+// The last key of the sorted buffered side, or `None` when it is empty or 
every key is NULL:
+// NULLs sort first, so the last key is NULL only when every buffered key is.
+fn buffered_extreme(values: &ArrayRef) -> Result<Option<ColumnarValue>> {
+    Ok(match values.len().checked_sub(1) {
+        // `apply_cmp` normalizes `-0.0` only in flat float scalars, not 
inside a
+        // `ScalarValue::Dictionary`, so normalize the key before taking it.
+        Some(last) if values.is_valid(last) => {
+            let extreme = normalize_float_zero(&values.slice(last, 1));
+            Some(ColumnarValue::Scalar(ScalarValue::try_from_array(
+                &extreme, 0,
+            )?))
+        }
+        _ => None,
+    })
+}
+
+// Which rows of `stream_values` can match at least one buffered row.
+//
+// Every match set is a suffix `[k, buffered_len)` of the sorted buffered 
side, so a streamed
+// row matches anything at all iff it matches the last buffered row, 
`buffered_extreme`
+// (`None` when no buffered key is non-null, so nothing can match). A NULL on 
either side
+// compares to NULL, which is no match.
+//
+// This must agree exactly with the scan's `JoinKeyComparator`, or it would 
change which rows
+// match. For flat keys one vectorized `apply_cmp` does: it normalizes `-0.0` 
to `+0.0` as the
+// comparator does. For nested keys it does not -- `apply_cmp` orders NULL 
elements inside a
+// key ascending, while the comparator applies the sort options (descending 
for `<`/`<=`) at
+// every level -- so those are decided with the scan's own comparator.
+fn matchable_rows(
+    stream_values: &ArrayRef,
+    buffered_values: &ArrayRef,
+    buffered_extreme: Option<&ColumnarValue>,
+    operator: Operator,
+    sort_options: SortOptions,
+) -> Result<BooleanArray> {
+    let num_rows = stream_values.len();
+    let Some(extreme) = buffered_extreme else {
+        return Ok(BooleanArray::new(BooleanBuffer::new_unset(num_rows), None));
+    };
+
+    // The comparator's nested ordering is itself wrong for `<`/`<=` (#25957). 
Once that is
+    // fixed, nested keys can go through `apply_cmp` too and this branch can 
be removed.
+    //
+    // Run-end encoded keys are decided here too: arrow's run-end comparison 
kernel overflows
+    // on an empty batch sliced past its first run, which the scan never 
reaches.
+    if stream_values.data_type().is_nested()
+        || matches!(stream_values.data_type(), DataType::RunEndEncoded(_, _))
+    {
+        let match_on_equal = matches_on_equal(operator)?;
+        let cmp = JoinKeyComparator::new(
+            &[Arc::clone(stream_values)],
+            &[Arc::clone(buffered_values)],
+            &[sort_options],
+            NullEquality::NullEqualsNothing,
+        )?;
+        let last = buffered_values.len() - 1;
+        let matchable = BooleanBuffer::collect_bool(num_rows, |row| {
+            stream_values.is_valid(row)

Review Comment:
   Addressed in 2627267b525c9eece069b89c55afc8099be349dc



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