jayzhan211 commented on code in PR #25584:
URL: https://github.com/apache/datafusion/pull/25584#discussion_r4173633387


##########
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##########
@@ -707,6 +721,127 @@ impl BitwiseSortMergeJoinStream {
         }
     }
 
+    /// Try initial probes or an exact summary, updating the matched bits.
+    /// Return true when this outer slice needs no ordinary filter evaluation.
+    /// Keep this work separate from the asynchronous spill and batch 
traversal.
+    #[inline(never)]
+    fn try_apply_summary(
+        &mut self,
+        outer_slice: &RecordBatch,
+        summary: &mut Option<Box<CachedSummary>>,
+        allow_summary: &mut Option<bool>,
+    ) -> Result<bool> {
+        let outer_group_start = self.outer_offset;
+        let outer_group_len = outer_slice.num_rows();
+        let filter = self.filter.as_ref().unwrap();
+
+        // Probe a small prefix before paying for min/max. This cost heuristic
+        // preserves early witnesses without adding an unconditional inner
+        // scan. A single outer row also stays on the ordinary path, which
+        // already visits each inner row at most once.
+        if summary.is_none()
+            && outer_group_len > 1
+            && *allow_summary.get_or_insert_with(|| {
+                self.inner_key_buffer
+                    .iter()
+                    .map(RecordBatch::num_rows)
+                    .sum::<usize>()
+                    >= 7
+            })
+        {
+            if let Some(comparison) = SemiAntiComparison::try_new(

Review Comment:
   `SemiAntiComparison::try_new` depends only on `filter`, `outer_is_left` and 
the two schemas, but it re-runs for every eligible key group until a summary is 
built. Recognize it once in `try_new` and keep it on the stream. That removes 
`CachedSummary`, the unsupported-filter arm of `allow_summary` and its lazy 
caching; eligibility becomes a plain `bool` per key group, and 
`try_apply_summary` reads `self.semi_anti_comparison.as_ref().unwrap()`.
   
   ```rs
       /// A residual that is one cross-side comparison, recognized once.
       /// `None` keeps ordinary filter evaluation for every key group.
       semi_anti_comparison: Option<SemiAntiComparison>,
   ```
   
   ```rs
           let is_mark = matches!(join_type, JoinType::LeftMark | 
JoinType::RightMark);
           let semi_anti_comparison = match &filter {
               Some(filter) if !is_mark => SemiAntiComparison::try_new(
                   filter,
                   outer_is_left,
                   outer.schema().as_ref(),
                   inner.schema().as_ref(),
               )?,
               _ => None,
           };
   ```
   
   ```diff
   -        let mut summary = None;
   -        // Delay the group-size check until an outer slice can use a 
summary,
   -        // then cache eligibility across batches. Spilled and mark joins 
never
   -        // qualify, regardless of their filter or group size.
   -        let mut allow_summary = (spill.is_some()
   -            || matches!(self.join_type, JoinType::LeftMark | 
JoinType::RightMark))
   -        .then_some(false);
   +        // Spilled groups keep ordinary evaluation.
   +        let can_summarize = spill.is_none()
   +            && self.semi_anti_comparison.is_some()
   +            && self
   +                .inner_key_buffer
   +                .iter()
   +                .map(RecordBatch::num_rows)
   +                .sum::<usize>()
   +                >= 7;
   +        let mut summary = None;
   ```



##########
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##########
@@ -707,6 +721,127 @@ impl BitwiseSortMergeJoinStream {
         }
     }
 
+    /// Try initial probes or an exact summary, updating the matched bits.
+    /// Return true when this outer slice needs no ordinary filter evaluation.
+    /// Keep this work separate from the asynchronous spill and batch 
traversal.
+    #[inline(never)]
+    fn try_apply_summary(
+        &mut self,
+        outer_slice: &RecordBatch,
+        summary: &mut Option<Box<CachedSummary>>,
+        allow_summary: &mut Option<bool>,
+    ) -> Result<bool> {
+        let outer_group_start = self.outer_offset;
+        let outer_group_len = outer_slice.num_rows();
+        let filter = self.filter.as_ref().unwrap();
+
+        // Probe a small prefix before paying for min/max. This cost heuristic
+        // preserves early witnesses without adding an unconditional inner
+        // scan. A single outer row also stays on the ordinary path, which
+        // already visits each inner row at most once.
+        if summary.is_none()
+            && outer_group_len > 1
+            && *allow_summary.get_or_insert_with(|| {
+                self.inner_key_buffer
+                    .iter()
+                    .map(RecordBatch::num_rows)
+                    .sum::<usize>()
+                    >= 7
+            })
+        {
+            if let Some(comparison) = SemiAntiComparison::try_new(
+                filter,
+                self.outer_is_left,
+                self.outer.schema().as_ref(),
+                self.inner.schema().as_ref(),
+            )? {
+                let mut remaining = 2;

Review Comment:
   With a fixed 2-row probe, a group whose witnesses start at inner row 2 pays 
a full min/max pass: O(N) where base stops after k evaluations. I added a case 
to the PR's harness (16 groups × 64 outer × 4096 inner, rows 0–1 miss, row 2 
matches all). Result: base 191 µs → head 700 µs on Utf8 (3.66×), 164 → 194 µs 
on Int64 (1.18×), base noise ≤0.8%. Since this is unconditional, the regression 
needs a bound. Scaling the probe budget with group size brings it back to 
0.96×/0.99× of base, while no-match stays ~100× faster (27 µs vs 2.9 ms) and 
small groups are unchanged:
   
   ```diff
   -                let mut remaining = 2;
   +                // Only pay for a full min/max pass after spending a probe
   +                // budget proportional to the group size.
   +                let inner_rows: usize = self
   +                    .inner_key_buffer
   +                    .iter()
   +                    .map(RecordBatch::num_rows)
   +                    .sum();
   +                let mut remaining = (inner_rows / 64).max(2);
   ```
   
   Please add the case to the bench:
   
   ```rs
           // Rows 0 and 1 miss; row 2 matches every outer row.
           (
               "not_equal_third_witness",
               16,
               64,
               4096,
               true,
               Operator::NotEq,
               JoinType::LeftSemi,
           ),
   ```
   
   ```diff
   -                            7 + if early_match {
   +                            7 + if name == "not_equal_third_witness" {
   +                                usize::from(row % rows_per_group >= 2)
   +                            } else if early_match {
   ```



##########
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##########
@@ -707,6 +721,127 @@ impl BitwiseSortMergeJoinStream {
         }
     }
 
+    /// Try initial probes or an exact summary, updating the matched bits.
+    /// Return true when this outer slice needs no ordinary filter evaluation.
+    /// Keep this work separate from the asynchronous spill and batch 
traversal.
+    #[inline(never)]
+    fn try_apply_summary(
+        &mut self,
+        outer_slice: &RecordBatch,
+        summary: &mut Option<Box<CachedSummary>>,
+        allow_summary: &mut Option<bool>,
+    ) -> Result<bool> {
+        let outer_group_start = self.outer_offset;
+        let outer_group_len = outer_slice.num_rows();
+        let filter = self.filter.as_ref().unwrap();
+
+        // Probe a small prefix before paying for min/max. This cost heuristic
+        // preserves early witnesses without adding an unconditional inner
+        // scan. A single outer row also stays on the ordinary path, which
+        // already visits each inner row at most once.
+        if summary.is_none()
+            && outer_group_len > 1
+            && *allow_summary.get_or_insert_with(|| {
+                self.inner_key_buffer
+                    .iter()
+                    .map(RecordBatch::num_rows)
+                    .sum::<usize>()
+                    >= 7
+            })
+        {
+            if let Some(comparison) = SemiAntiComparison::try_new(
+                filter,
+                self.outer_is_left,
+                self.outer.schema().as_ref(),
+                self.inner.schema().as_ref(),
+            )? {
+                let mut remaining = 2;
+                // A NULL outer operand cannot match any inner value.
+                let matchable = outer_group_len
+                    - outer_slice.column(comparison.outer_column).null_count();
+                if matchable == 0 {
+                    return Ok(true);
+                }
+                for inner_batch in &self.inner_key_buffer {
+                    let len = remaining.min(inner_batch.num_rows());
+                    for row in 0..len {
+                        let result = comparison.evaluate_inner_row(
+                            inner_batch,
+                            row,
+                            outer_slice,
+                        )?;
+                        let values = result.values();
+                        apply_bitwise_binary_op(
+                            self.matched.as_slice_mut(),
+                            outer_group_start,
+                            values.inner().as_slice(),
+                            values.offset(),
+                            outer_group_len,
+                            |a, b| a | b,
+                        );
+                        let matched_count = UnalignedBitChunk::new(
+                            self.matched.as_slice(),
+                            outer_group_start,
+                            outer_group_len,
+                        )
+                        .count_ones();
+                        if matched_count == matchable {
+                            return Ok(true);
+                        }
+                    }
+                    remaining -= len;
+                    if remaining == 0 {
+                        break;
+                    }
+                }
+                let bytes = comparison

Review Comment:
   The extrema are copies of two values already covered by the 
`inner_key_buffer` reservation, so admitting `6 ×` the largest parent string 
buffer before reducing isn't needed. Today each summarized group pays a 
`try_grow` + `shrink` (the FairSpillPool lock). It also briefly over-reserves 
enough to push concurrent consumers into spilling, and inflates 
`peak_mem_used`. Summarize first, then shrink to the summary's size; this 
removes both `memory_size` fns and the refusal branch. Prototype: 
`small_group_seven_no_match` 0.97×, and 0.88× with a shared FairSpillPool 
across 4 tasks (that case is noisy, ±8%).
   
   ```diff
   -                let bytes = comparison
   -                    .memory_size(&self.inner_key_buffer)
   -                    .saturating_add(size_of::<SemiAntiComparison>());
   -                if self.reservation.try_grow(bytes).is_ok() {
   -                    self.peak_mem_used.set_max(self.reservation.size());
   -                    // Keep scalar state out of the surrounding async 
futures.
   -                    // Admit the comparison descriptor alongside its 
scalars.
   -                    let values = 
comparison.summarize(&self.inner_key_buffer)?;
   -                    let retained = values
   -                        .memory_size()
   -                        .saturating_add(size_of::<SemiAntiComparison>());
   -                    debug_assert!(retained <= bytes);
   -                    let cached = Box::new(CachedSummary { comparison, 
values });
   -                    // Extrema own their values. Once reduction succeeds, 
the
   -                    // original group is no longer needed, including for 
outer
   -                    // rows arriving in later batches. Release it before any
   -                    // further child awaits and keep only the probe 
allowance.
   -                    self.inner_key_buffer.clear();
   -                    self.inner_buffer_size = 0;
   -                    self.reservation.shrink(self.reservation.size() - 
retained);
   -                    *summary = Some(cached);
   -                } else {
   -                    // Input is still buffered, so refused summary admission
   -                    // does not add a memory error or disable existing 
spilling.
   -                    *allow_summary = Some(false);
   -                }
   +                // The extrema copy two buffered values, so they fit in the
   +                // reservation already held for `inner_key_buffer`.
   +                let values = comparison.summarize(&self.inner_key_buffer)?;
   +                self.inner_key_buffer.clear();
   +                self.inner_buffer_size = 0;
   +                self.reservation.resize(values.size());
   +                *summary = Some(Box::new(CachedSummary { comparison, values 
}));
   ```
   
   `GroupSummary::memory_size` becomes:
   
   ```rs
   impl GroupSummary {
       pub(super) fn size(&self) -> usize {
           self.extreme.size() + self.other.as_ref().map_or(0, 
ScalarValue::size)
       }
   }
   ```
   
   In `dropping_stream_releases_buffered_group_and_summary`, the peak no longer 
includes an admitted overlap. The middle-budget case in 
`large_strings_fall_back_when_summary_or_buffer_exceeds_memory` no longer 
exercises anything and can go.
   
   ```diff
   -        assert!(
   -            peak > inner.get_array_memory_size(),
   -            "summary storage must be admitted alongside the buffered group"
   -        );
   +        assert_eq!(peak, inner.get_array_memory_size());
   ```



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