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


##########
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##########
@@ -912,6 +1003,122 @@ impl BitwiseSortMergeJoinStream {
         }
     }
 
+    /// Return a bounded slice length, yielding periodically even if every
+    /// child batch is immediately ready. Pausing the timer excludes scheduling
+    /// time from the join's own work.
+    async fn summary_slice_len(&mut self, remaining: usize) -> usize {

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. Removed the explicit 
yields and 1024-row slicing. Reduction now uses the existing buffered group and 
whole-batch min/max accumulator updates; the existing input/output await points 
remain. The Rust drop test checks admitted memory release rather than a yield 
count.



##########
datafusion/physical-plan/src/joins/sort_merge_join/existence_summary.rs:
##########
@@ -0,0 +1,688 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Exact bounded summaries of an ordinary semi/anti join's residual predicate.
+//!
+//! For one equi-key group, existence of `outer != inner` needs at most one
+//! non-null representative and a second-distinct-value bit. Existence of
+//! `outer <[=] inner` needs the inner maximum, and `outer >[=] inner` needs 
its
+//! minimum. Existence distributes over OR, but not AND: two comparisons may
+//! have different witnesses. Each compiled disjunct therefore has at most one
+//! cross-side comparison, with any total side-local guards attached to it.
+//!
+//! Results are non-null *existence* bits, not the residual's three-valued SQL
+//! result. A null comparison is not a witness. The caller supplies only rows
+//! from one matching equi-key group and implements semi/anti polarity itself.
+//! Mark and null-aware joins must not use this interface.
+//!
+//! Input batches are borrowed only during update/evaluation. A representative
+//! is copied after reserving its storage, never retained as a slice of a 
source
+//! array. The same dedicated reservation must accompany updates and resets.
+//! Memory admission failure is an execution error; after input has been
+//! consumed it is not safe to fall back without an explicit replay path.
+
+use std::cmp::Ordering;
+use std::sync::Arc;
+
+use arrow::array::{Array, ArrayRef, AsArray, BooleanArray, RecordBatch};
+use arrow::compute::SortOptions;
+use arrow::compute::kernels::cmp::{gt, gt_eq, lt, lt_eq, neq};
+use arrow::datatypes::{DataType, Schema};
+use arrow_ord::ord::make_comparator;
+use datafusion_common::{JoinSide, Result, ScalarValue};
+use datafusion_execution::memory_pool::MemoryReservation;
+use datafusion_expr::Operator;
+use datafusion_physical_expr::PhysicalExpr;
+use datafusion_physical_expr::expressions::{
+    BinaryExpr, CaseExpr, Column, IsNotNullExpr, IsNullExpr, Literal, NotExpr,
+};
+
+use crate::joins::utils::JoinFilter;
+
+type Expr = Arc<dyn PhysicalExpr>;
+
+// Bound compilation work, expression recursion and the number of retained
+// representatives even if distributive normalization would grow exponentially.
+const MAX_CLAUSES: usize = 32;
+const MAX_DEPTH: usize = 64;
+
+#[derive(Debug, Clone)]
+struct Comparison {
+    outer: Expr,
+    inner: Expr,
+    op: Operator,
+}
+
+#[derive(Debug, Default)]
+struct State {

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. Replaced the custom 
reduction with `MinAccumulator` / `MaxAccumulator` from aggregate-common. `<>` 
uses `x <> min OR x <> max`, without representative/`multiple` state or the 
per-row scalar reduction loop. Because only direct columns qualify now, the 
helper calls those accumulators directly without moving the expression/guard 
wrapper from hash join.



##########
datafusion/physical-plan/src/joins/sort_merge_join/existence_summary.rs:
##########
@@ -0,0 +1,688 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Exact bounded summaries of an ordinary semi/anti join's residual predicate.
+//!
+//! For one equi-key group, existence of `outer != inner` needs at most one
+//! non-null representative and a second-distinct-value bit. Existence of
+//! `outer <[=] inner` needs the inner maximum, and `outer >[=] inner` needs 
its
+//! minimum. Existence distributes over OR, but not AND: two comparisons may
+//! have different witnesses. Each compiled disjunct therefore has at most one
+//! cross-side comparison, with any total side-local guards attached to it.
+//!
+//! Results are non-null *existence* bits, not the residual's three-valued SQL
+//! result. A null comparison is not a witness. The caller supplies only rows
+//! from one matching equi-key group and implements semi/anti polarity itself.
+//! Mark and null-aware joins must not use this interface.
+//!
+//! Input batches are borrowed only during update/evaluation. A representative
+//! is copied after reserving its storage, never retained as a slice of a 
source
+//! array. The same dedicated reservation must accompany updates and resets.
+//! Memory admission failure is an execution error; after input has been
+//! consumed it is not safe to fall back without an explicit replay path.
+
+use std::cmp::Ordering;
+use std::sync::Arc;
+
+use arrow::array::{Array, ArrayRef, AsArray, BooleanArray, RecordBatch};
+use arrow::compute::SortOptions;
+use arrow::compute::kernels::cmp::{gt, gt_eq, lt, lt_eq, neq};
+use arrow::datatypes::{DataType, Schema};
+use arrow_ord::ord::make_comparator;
+use datafusion_common::{JoinSide, Result, ScalarValue};
+use datafusion_execution::memory_pool::MemoryReservation;
+use datafusion_expr::Operator;
+use datafusion_physical_expr::PhysicalExpr;
+use datafusion_physical_expr::expressions::{
+    BinaryExpr, CaseExpr, Column, IsNotNullExpr, IsNullExpr, Literal, NotExpr,
+};
+
+use crate::joins::utils::JoinFilter;
+
+type Expr = Arc<dyn PhysicalExpr>;
+
+// Bound compilation work, expression recursion and the number of retained
+// representatives even if distributive normalization would grow exponentially.
+const MAX_CLAUSES: usize = 32;
+const MAX_DEPTH: usize = 64;
+
+#[derive(Debug, Clone)]
+struct Comparison {
+    outer: Expr,
+    inner: Expr,
+    op: Operator,
+}
+
+#[derive(Debug, Default)]
+struct State {
+    /// Used for clauses with only side-local predicates. Even an outer TRUE
+    /// needs at least one qualifying inner row to witness existence.
+    present: bool,
+    representative: Option<ScalarValue>,
+    multiple: bool,
+    reserved: usize,
+}
+
+impl State {
+    fn set_multiple(&mut self, reservation: &MemoryReservation) {
+        self.representative = None;
+        reservation.shrink(self.reserved);
+        *self = Self {
+            multiple: true,
+            ..Self::default()
+        };
+    }
+
+    /// Admit the simultaneous old and new copies before copying the candidate.
+    /// For ranges, retain it only if it improves the current extremum.
+    fn update_from_array(
+        &mut self,
+        array: &ArrayRef,
+        index: usize,
+        order: Option<Ordering>,
+        reservation: &MemoryReservation,
+        peak: &mut usize,
+    ) -> Result<()> {
+        let bytes = scalar_storage_size(array, index);
+        reservation.try_grow(bytes)?;
+        *peak = (*peak).max(reservation.size());
+        let candidate = ScalarValue::try_from_array(array, index)?;
+        if let Some(previous) = &self.representative
+            && order.is_some_and(|order| candidate.partial_cmp(previous) != 
Some(order))
+        {
+            drop(candidate);
+            reservation.shrink(bytes);
+            return Ok(());
+        }
+        self.representative = Some(candidate);
+        reservation.shrink(self.reserved);
+        self.reserved = bytes;
+        Ok(())
+    }
+}
+
+#[derive(Debug, Default)]
+struct Clause {
+    outer_guard: Option<Expr>,
+    inner_guard: Option<Expr>,
+    comparison: Option<Comparison>,
+    state: State,
+}
+
+/// A compiled residual and the bounded state for its current inner key group.
+#[derive(Debug)]
+pub(super) struct ExistenceSummary {
+    clauses: Vec<Clause>,
+    peak: usize,
+}
+
+impl ExistenceSummary {
+    /// Compile only exact, total expressions over supported identically typed
+    /// values. `None` means the ordinary residual path must be retained.
+    pub(super) fn try_new(
+        filter: &JoinFilter,
+        outer_is_left: bool,
+        outer_schema: &Schema,
+        inner_schema: &Schema,
+    ) -> Result<Option<Self>> {
+        let compiler = Compiler {
+            filter,
+            outer_is_left,
+            outer_schema,
+            inner_schema,
+        };
+        Ok(compiler
+            .compile(filter.expression(), 0)?
+            .map(|clauses| Self { clauses, peak: 0 }))
+    }
+
+    /// Incorporate one slice of the current inner group. The caller bounds
+    /// slice size for cancellation latency and continues draining the group
+    /// even when every not-equal summary has saturated.
+    pub(super) fn update(
+        &mut self,
+        inner: &RecordBatch,
+        reservation: &MemoryReservation,
+    ) -> Result<()> {
+        for clause in &mut self.clauses {
+            let guard = evaluate_guard(clause.inner_guard.as_ref(), inner)?;
+            let selected = |row| {
+                guard
+                    .as_ref()
+                    .is_none_or(|g| g.is_valid(row) && g.value(row))
+            };
+            let Some(comparison) = &clause.comparison else {
+                clause.state.present |= (0..inner.num_rows()).any(selected);
+                continue;
+            };
+            if clause.state.multiple {
+                continue;
+            }
+            let values = comparison
+                .inner
+                .evaluate(inner)?
+                .into_array(inner.num_rows())?;
+            let Some(first) =
+                (0..values.len()).find(|&row| selected(row) && 
values.is_valid(row))
+            else {
+                continue;
+            };
+            if comparison.op == Operator::NotEq {
+                if clause.state.representative.is_none() {
+                    clause.state.update_from_array(
+                        &values,
+                        first,
+                        None,
+                        reservation,
+                        &mut self.peak,
+                    )?;
+                }
+                let value = clause.state.representative.as_ref().unwrap();
+                for row in first..values.len() {
+                    if selected(row)
+                        && values.is_valid(row)
+                        && !value.eq_array(&values, row)?
+                    {
+                        clause.state.set_multiple(reservation);
+                        break;
+                    }
+                }
+            } else {
+                let order = if matches!(comparison.op, Operator::Lt | 
Operator::LtEq) {
+                    Ordering::Greater
+                } else {
+                    Ordering::Less
+                };
+                let comparator = make_comparator(
+                    values.as_ref(),
+                    values.as_ref(),
+                    SortOptions::default(),
+                )?;
+                let mut candidate = first;
+                for row in first + 1..values.len() {
+                    if selected(row)
+                        && values.is_valid(row)
+                        && comparator(row, candidate) == order
+                    {
+                        candidate = row;
+                    }
+                }
+                clause.state.update_from_array(
+                    &values,
+                    candidate,
+                    Some(order),
+                    reservation,
+                    &mut self.peak,
+                )?;
+            }
+        }
+        Ok(())
+    }
+
+    /// Return non-null witness bits for an outer slice. Null outer values do
+    /// not match even when a not-equal summary contains two distinct values.
+    pub(super) fn evaluate(&self, outer: &RecordBatch) -> Result<BooleanArray> 
{
+        let mut matches = vec![false; outer.num_rows()];

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. The probe uses 
`apply_cmp` and `boolean_mask_from_filter`, then merges the result directly 
into the existing match bitmap. `<>` ORs the two extremum comparisons. Compound 
predicates use the existing filter evaluator, so the row-by-clause loop, 
boolean repack, and separate comparison table are gone.



##########
datafusion/physical-plan/src/joins/sort_merge_join/bitwise_stream.rs:
##########
@@ -912,6 +1003,122 @@ impl BitwiseSortMergeJoinStream {
         }
     }
 
+    /// Return a bounded slice length, yielding periodically even if every
+    /// child batch is immediately ready. Pausing the timer excludes scheduling
+    /// time from the join's own work.
+    async fn summary_slice_len(&mut self, remaining: usize) -> usize {
+        if self.existence_summary.as_ref().unwrap().rows_until_yield == 0 {
+            self.stop_join_time();
+            tokio::task::yield_now().await;
+            self.start_join_time();
+            self.existence_summary.as_mut().unwrap().rows_until_yield =
+                SUMMARY_WORK_BUDGET;
+        }
+        let state = self.existence_summary.as_mut().unwrap();
+        let len = remaining.min(state.rows_until_yield);
+        state.rows_until_yield -= len;
+        len
+    }
+
+    /// Consume the whole inner key group into owned summary state. Continue
+    /// draining after saturation so input errors and group boundaries retain
+    /// their ordinary behavior.
+    async fn summarize_inner_key_group(&mut self) -> Result<()> {

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. Removed both parallel 
traversal loops. This version reuses `buffer_inner_key_group` and the filtered 
match loop, caches one summary per key across outer batches, and releases the 
buffered group after successful reduction. The cached scalar state is boxed, 
and initial probes borrow input batches without allocating extra slices. 
Group-size eligibility is checked only when a multirow outer slice can use a 
summary, then cached across batches. Spilled groups retain ordinary evaluation; 
operands are direct column indices. Matched base/head measurements exposed 
fixed-cost regressions on small groups, so reduction now requires at least 
seven inner rows and more than one row in the current outer slice, with up to 
two initial row comparisons first. The harness now covers three-, four-, and 
seven-row groups with single, shared-pool, and separate-pool execution. The 
traversal cleanup is complete, but I am leaving this thread open for the re
 maining common-path performance discussion. Performance review remains open. 
The unfiltered semi timing difference is strongly sensitive to benchmark binary 
layout: it reproduced with identical production code and identical timed work, 
while native profiles locate the difference in unchanged Arrow key-comparison 
code. The adverse single-row range result did not reproduce under native 
sampling, which does not establish that it is fixed. SQL Q13 also has an 
adverse mean with substantial process variation. I have retained all raw 
results and diagnostic provenance, and am not claiming blanket non-regression 
or performance readiness.



##########
datafusion/sqllogictest/test_files/sort_merge_join.slt:
##########
@@ -971,6 +971,47 @@ WHERE t1_sorted.data < 0
 ----
 100
 
+statement ok
+SET datafusion.execution.enable_sort_merge_join_existence_summary = true;
+
+# OR clauses can have different witnesses. A null component is not itself
+# a conflict, and a missing inner group must survive an anti join.
+statement ok
+CREATE TABLE identities(id BIGINT, k BIGINT, user_id VARCHAR, org_id VARCHAR) 
AS VALUES
+  (1, 1, 'alice', 'x'), (2, 2, 'alice', 'x'),
+  (3, 3, NULL, 'x'), (4, 4, 'alice', NULL), (5, 5, 'carol', 'x');
+
+statement ok
+CREATE TABLE observations(k BIGINT, user_id VARCHAR, org_id VARCHAR, active 
BOOLEAN) AS VALUES
+  (1, 'alice', 'x', true), (1, 'alice', 'x', true),
+  (1, 'bob', 'y', false), (2, 'alice', 'y', true),
+  (2, 'bob', 'x', true), (3, 'bob', NULL, true), (4, 'bob', NULL, true);
+
+query I rowsort
+SELECT l.id FROM identities l LEFT ANTI JOIN observations r
+ON l.k = r.k AND r.active AND (l.user_id <> r.user_id OR l.org_id <> r.org_id);

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. Agreed: the top-level 
conjunct did not establish that an inner guard reached the summary. Guard 
compilation is removed in this narrower version. The SLT now includes a guard 
inside one OR arm for generic-path coverage; physical-plan inspection confirms 
it remains in the join residual. The Rust recognizer test also supplies 
guarded/compound physical predicates directly and asserts rejection, so SQL 
pushdown cannot hide an eligibility mistake.



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