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


##########
datafusion/physical-plan/src/joins/sort_merge_join/existence_summary_tests.rs:
##########
@@ -0,0 +1,753 @@
+// 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.
+
+//! Differential execution tests for bounded semi/anti existence summaries.
+
+use std::sync::Arc;
+use std::task::Context;
+
+use arrow::array::{
+    Array, ArrayRef, BooleanArray, Date32Array, Decimal128Array, Int32Array,
+    LargeStringArray, RecordBatch, StringArray, StringViewArray,
+    TimestampMicrosecondArray,
+};
+use arrow::compute::SortOptions;
+use arrow::datatypes::{DataType, Field, Schema};
+use datafusion_common::{
+    DataFusionError, JoinSide, JoinType, NullEquality, Result, ScalarValue,
+    assert_contains,
+};
+use datafusion_execution::TaskContext;
+use datafusion_execution::config::SessionConfig;
+use datafusion_execution::runtime_env::RuntimeEnvBuilder;
+use datafusion_expr::Operator;
+use datafusion_physical_expr::expressions::{
+    BinaryExpr, CaseExpr, Column, IsNullExpr, Literal, NotExpr,
+};
+
+use crate::joins::SortMergeJoinExec;
+use crate::joins::utils::{ColumnIndex, JoinFilter};
+use crate::test::TestMemoryExec;
+use crate::test::exec::MockExec;
+use crate::{ExecutionPlan, PhysicalExpr, common};
+
+const JOINS: [JoinType; 4] = [
+    JoinType::LeftSemi,
+    JoinType::LeftAnti,
+    JoinType::RightSemi,
+    JoinType::RightAnti,
+];
+type Row = (Option<i32>, Option<i32>, Option<bool>);
+
+fn column(index: usize) -> Arc<dyn PhysicalExpr> {
+    Arc::new(Column::new(
+        ["left_value", "right_value", "left_guard", "right_guard"][index],
+        index,
+    ))
+}
+
+fn binary(
+    left: Arc<dyn PhysicalExpr>,
+    op: Operator,
+    right: Arc<dyn PhysicalExpr>,
+) -> Arc<dyn PhysicalExpr> {
+    Arc::new(BinaryExpr::new(left, op, right))
+}
+
+fn filter(expr: Arc<dyn PhysicalExpr>, value_type: DataType) -> JoinFilter {
+    JoinFilter::new(
+        expr,
+        vec![
+            ColumnIndex {
+                index: 1,
+                side: JoinSide::Left,
+            },
+            ColumnIndex {
+                index: 1,
+                side: JoinSide::Right,
+            },
+            ColumnIndex {
+                index: 2,
+                side: JoinSide::Left,
+            },
+            ColumnIndex {
+                index: 2,
+                side: JoinSide::Right,
+            },
+        ],
+        Arc::new(Schema::new(vec![
+            Field::new("left_value", value_type.clone(), true),
+            Field::new("right_value", value_type, true),
+            Field::new("left_guard", DataType::Boolean, true),
+            Field::new("right_guard", DataType::Boolean, true),
+        ])),
+    )
+}
+
+fn comparison(op: Operator) -> JoinFilter {
+    filter(binary(column(0), op, column(1)), DataType::Int32)
+}
+
+fn batch(rows: &[Row]) -> Result<RecordBatch> {
+    batch_values(
+        rows,
+        Arc::new(Int32Array::from_iter(rows.iter().map(|r| r.1))),
+    )
+}
+
+fn batch_values(rows: &[Row], values: ArrayRef) -> Result<RecordBatch> {
+    RecordBatch::try_from_iter(vec![
+        (
+            "key",
+            Arc::new(Int32Array::from_iter(rows.iter().map(|r| r.0))) as 
ArrayRef,
+        ),
+        ("value", values),
+        (
+            "guard",
+            Arc::new(BooleanArray::from_iter(rows.iter().map(|r| r.2))),
+        ),
+        (
+            "id",
+            Arc::new(Int32Array::from_iter_values(0..rows.len() as i32)),
+        ),
+    ])
+    .map_err(Into::into)
+}
+
+fn input(batch: &RecordBatch, chunk: usize) -> Result<Arc<dyn ExecutionPlan>> {
+    // Empty batches between slices exercise both empty input and non-zero 
array offsets.
+    let mut batches = vec![batch.slice(0, 0)];
+    for offset in (0..batch.num_rows()).step_by(chunk) {
+        batches.push(batch.slice(offset, chunk.min(batch.num_rows() - 
offset)));
+        batches.push(batch.slice(offset, 0));
+    }
+    Ok(TestMemoryExec::try_new_exec(
+        &[batches],
+        batch.schema(),
+        None,
+    )?)
+}
+
+fn join(
+    left: Arc<dyn ExecutionPlan>,
+    right: Arc<dyn ExecutionPlan>,
+    join_type: JoinType,
+    filter: JoinFilter,
+    options: SortOptions,
+    nulls: NullEquality,
+) -> Result<SortMergeJoinExec> {
+    SortMergeJoinExec::try_new(
+        left,
+        right,
+        vec![(
+            Arc::new(Column::new("key", 0)),
+            Arc::new(Column::new("key", 0)),
+        )],
+        Some(filter),
+        join_type,
+        vec![options],
+        nulls,
+    )
+}
+
+fn config(batch_size: usize, enabled: bool) -> SessionConfig {
+    let mut config = SessionConfig::new().with_batch_size(batch_size);
+    config
+        .options_mut()
+        .execution
+        .enable_sort_merge_join_existence_summary = enabled;
+    config
+}
+
+fn context(batch_size: usize, enabled: bool) -> Arc<TaskContext> {
+    Arc::new(TaskContext::default().with_session_config(config(batch_size, 
enabled)))
+}
+
+fn metric(plan: &dyn ExecutionPlan, name: &str) -> usize {
+    plan.metrics()
+        .unwrap()
+        .iter()
+        .filter(|metric| metric.value().name() == name)
+        .map(|metric| metric.value().as_usize())
+        .sum()
+}
+
+async fn ids(plan: &dyn ExecutionPlan, ctx: Arc<TaskContext>) -> 
Result<Vec<i32>> {
+    let mut ids = common::collect(plan.execute(0, ctx)?)
+        .await?
+        .iter()
+        .flat_map(|batch| {
+            batch
+                .column(3)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap()
+                .values()
+                .to_vec()
+        })
+        .collect::<Vec<_>>();
+    ids.sort_unstable();
+    Ok(ids)
+}
+
+#[tokio::test]
+async fn comparisons_match_generic_execution_across_groups_and_orientations() 
-> Result<()>

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. Moved result/type/batch 
coverage into `sort_merge_join_matrix.slt`: all five comparisons and four 
semi/anti orientations, NULLs, duplicates, three string layouts, 
temporal/decimal values, and expression fallback. Seven-row fixtures keep the 
summary reachable at its cost threshold; the existing hash-join/batch-size 
matrix still covers ordinary execution too. Rust now concentrates on physical 
NOT/eligibility, memory admission/spilling, input errors, and dropping active 
state. The flag, five counters, and core toggle test are removed; removing the 
flag makes the optimization available under default configuration. The existing 
randomized fuzz inputs are sparse, so seven-row SQL fixtures cover eligible 
comparisons at the cost threshold, while Rust reservation assertions prove 
summary construction and reuse across one-, two-, then one-row outer slices. 
The final source passed 12,211 Rust tests (including all 44 
default-configuration j
 oin-fuzz tests), all 525 SQL logic files, all-target/all-feature Clippy, 
formatting, and the full lint suite.



##########
datafusion/physical-plan/src/joins/sort_merge_join/existence_summary_tests.rs:
##########
@@ -0,0 +1,753 @@
+// 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.
+
+//! Differential execution tests for bounded semi/anti existence summaries.
+
+use std::sync::Arc;
+use std::task::Context;
+
+use arrow::array::{
+    Array, ArrayRef, BooleanArray, Date32Array, Decimal128Array, Int32Array,
+    LargeStringArray, RecordBatch, StringArray, StringViewArray,
+    TimestampMicrosecondArray,
+};
+use arrow::compute::SortOptions;
+use arrow::datatypes::{DataType, Field, Schema};
+use datafusion_common::{
+    DataFusionError, JoinSide, JoinType, NullEquality, Result, ScalarValue,
+    assert_contains,
+};
+use datafusion_execution::TaskContext;
+use datafusion_execution::config::SessionConfig;
+use datafusion_execution::runtime_env::RuntimeEnvBuilder;
+use datafusion_expr::Operator;
+use datafusion_physical_expr::expressions::{
+    BinaryExpr, CaseExpr, Column, IsNullExpr, Literal, NotExpr,
+};
+
+use crate::joins::SortMergeJoinExec;
+use crate::joins::utils::{ColumnIndex, JoinFilter};
+use crate::test::TestMemoryExec;
+use crate::test::exec::MockExec;
+use crate::{ExecutionPlan, PhysicalExpr, common};
+
+const JOINS: [JoinType; 4] = [
+    JoinType::LeftSemi,
+    JoinType::LeftAnti,
+    JoinType::RightSemi,
+    JoinType::RightAnti,
+];
+type Row = (Option<i32>, Option<i32>, Option<bool>);
+
+fn column(index: usize) -> Arc<dyn PhysicalExpr> {
+    Arc::new(Column::new(
+        ["left_value", "right_value", "left_guard", "right_guard"][index],
+        index,
+    ))
+}
+
+fn binary(
+    left: Arc<dyn PhysicalExpr>,
+    op: Operator,
+    right: Arc<dyn PhysicalExpr>,
+) -> Arc<dyn PhysicalExpr> {
+    Arc::new(BinaryExpr::new(left, op, right))
+}
+
+fn filter(expr: Arc<dyn PhysicalExpr>, value_type: DataType) -> JoinFilter {
+    JoinFilter::new(
+        expr,
+        vec![
+            ColumnIndex {
+                index: 1,
+                side: JoinSide::Left,
+            },
+            ColumnIndex {
+                index: 1,
+                side: JoinSide::Right,
+            },
+            ColumnIndex {
+                index: 2,
+                side: JoinSide::Left,
+            },
+            ColumnIndex {
+                index: 2,
+                side: JoinSide::Right,
+            },
+        ],
+        Arc::new(Schema::new(vec![
+            Field::new("left_value", value_type.clone(), true),
+            Field::new("right_value", value_type, true),
+            Field::new("left_guard", DataType::Boolean, true),
+            Field::new("right_guard", DataType::Boolean, true),
+        ])),
+    )
+}
+
+fn comparison(op: Operator) -> JoinFilter {
+    filter(binary(column(0), op, column(1)), DataType::Int32)
+}
+
+fn batch(rows: &[Row]) -> Result<RecordBatch> {
+    batch_values(
+        rows,
+        Arc::new(Int32Array::from_iter(rows.iter().map(|r| r.1))),
+    )
+}
+
+fn batch_values(rows: &[Row], values: ArrayRef) -> Result<RecordBatch> {
+    RecordBatch::try_from_iter(vec![
+        (
+            "key",
+            Arc::new(Int32Array::from_iter(rows.iter().map(|r| r.0))) as 
ArrayRef,
+        ),
+        ("value", values),
+        (
+            "guard",
+            Arc::new(BooleanArray::from_iter(rows.iter().map(|r| r.2))),
+        ),
+        (
+            "id",
+            Arc::new(Int32Array::from_iter_values(0..rows.len() as i32)),
+        ),
+    ])
+    .map_err(Into::into)
+}
+
+fn input(batch: &RecordBatch, chunk: usize) -> Result<Arc<dyn ExecutionPlan>> {
+    // Empty batches between slices exercise both empty input and non-zero 
array offsets.
+    let mut batches = vec![batch.slice(0, 0)];
+    for offset in (0..batch.num_rows()).step_by(chunk) {
+        batches.push(batch.slice(offset, chunk.min(batch.num_rows() - 
offset)));
+        batches.push(batch.slice(offset, 0));
+    }
+    Ok(TestMemoryExec::try_new_exec(
+        &[batches],
+        batch.schema(),
+        None,
+    )?)
+}
+
+fn join(
+    left: Arc<dyn ExecutionPlan>,
+    right: Arc<dyn ExecutionPlan>,
+    join_type: JoinType,
+    filter: JoinFilter,
+    options: SortOptions,
+    nulls: NullEquality,
+) -> Result<SortMergeJoinExec> {
+    SortMergeJoinExec::try_new(
+        left,
+        right,
+        vec![(
+            Arc::new(Column::new("key", 0)),
+            Arc::new(Column::new("key", 0)),
+        )],
+        Some(filter),
+        join_type,
+        vec![options],
+        nulls,
+    )
+}
+
+fn config(batch_size: usize, enabled: bool) -> SessionConfig {
+    let mut config = SessionConfig::new().with_batch_size(batch_size);
+    config
+        .options_mut()
+        .execution
+        .enable_sort_merge_join_existence_summary = enabled;
+    config
+}
+
+fn context(batch_size: usize, enabled: bool) -> Arc<TaskContext> {
+    Arc::new(TaskContext::default().with_session_config(config(batch_size, 
enabled)))
+}
+
+fn metric(plan: &dyn ExecutionPlan, name: &str) -> usize {
+    plan.metrics()
+        .unwrap()
+        .iter()
+        .filter(|metric| metric.value().name() == name)
+        .map(|metric| metric.value().as_usize())
+        .sum()
+}
+
+async fn ids(plan: &dyn ExecutionPlan, ctx: Arc<TaskContext>) -> 
Result<Vec<i32>> {
+    let mut ids = common::collect(plan.execute(0, ctx)?)
+        .await?
+        .iter()
+        .flat_map(|batch| {
+            batch
+                .column(3)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap()
+                .values()
+                .to_vec()
+        })
+        .collect::<Vec<_>>();
+    ids.sort_unstable();
+    Ok(ids)
+}
+
+#[tokio::test]
+async fn comparisons_match_generic_execution_across_groups_and_orientations() 
-> Result<()>
+{
+    let left = vec![
+        (None, None, None),
+        (None, Some(4), Some(true)),
+        (Some(0), Some(7), Some(true)),
+        (Some(1), None, None),
+        (Some(1), Some(2), Some(true)),
+        (Some(1), Some(2), Some(false)),
+        (Some(1), Some(6), None),
+        (Some(2), Some(3), Some(true)),
+        (Some(4), Some(5), Some(false)),
+    ];
+    let right = vec![
+        (None, Some(1), Some(true)),
+        (Some(1), None, None),
+        (Some(1), Some(2), Some(false)),
+        (Some(1), Some(2), Some(true)),
+        (Some(1), Some(6), Some(true)),
+        (Some(2), None, None),
+        (Some(3), Some(4), Some(false)),
+    ];
+    let mut cases = vec![];
+    for op in [
+        Operator::NotEq,
+        Operator::Lt,
+        Operator::LtEq,
+        Operator::Gt,
+        Operator::GtEq,
+    ] {
+        for kind in JOINS {
+            cases.push((
+                op,
+                kind,
+                SortOptions::default(),
+                NullEquality::NullEqualsNothing,
+            ));
+        }
+    }
+    // Key ordering and null equality are independent of residual comparison.
+    cases.extend([
+        (
+            Operator::NotEq,
+            JoinType::LeftAnti,
+            SortOptions {
+                descending: false,
+                nulls_first: false,
+            },
+            NullEquality::NullEqualsNull,
+        ),
+        (
+            Operator::Lt,
+            JoinType::RightSemi,
+            SortOptions {
+                descending: true,
+                nulls_first: true,
+            },
+            NullEquality::NullEqualsNull,
+        ),
+        (
+            Operator::LtEq,
+            JoinType::LeftSemi,
+            SortOptions {
+                descending: true,
+                nulls_first: false,
+            },
+            NullEquality::NullEqualsNothing,
+        ),
+    ]);
+    for (op, kind, options, nulls) in cases {
+        let sort = |rows: &[Row]| {
+            let mut rows = rows.to_vec();
+            rows.sort_by(|a, b| {
+                let null_order = if options.nulls_first {
+                    b.0.is_none().cmp(&a.0.is_none())
+                } else {
+                    a.0.is_none().cmp(&b.0.is_none())
+                };
+                null_order.then_with(|| {
+                    if options.descending {
+                        b.0.cmp(&a.0)
+                    } else {
+                        a.0.cmp(&b.0)
+                    }
+                })
+            });
+            rows
+        };
+        let (left, right) = (batch(&sort(&left))?, batch(&sort(&right))?);
+        let mut outputs = vec![];
+        for enabled in [false, true] {
+            let plan = join(
+                input(&left, 2)?,
+                input(&right, 3)?,
+                kind,
+                comparison(op),
+                options,
+                nulls,
+            )?;
+            outputs.push(ids(&plan, context(2, enabled)).await?);
+            assert_eq!(
+                metric(&plan, "existence_summary_enabled"),
+                usize::from(enabled)
+            );
+            if enabled {
+                assert!(metric(&plan, "existence_summary_inner_rows") > 0);
+            } else {
+                assert!(plan.metrics().unwrap().iter().all(|metric| {
+                    !metric.value().name().starts_with("existence_summary_")
+                }));
+            }
+        }
+        assert_eq!(
+            outputs[0], outputs[1],
+            "{kind:?} {op:?} {options:?} {nulls:?}"
+        );
+        if kind == JoinType::LeftSemi
+            && options == SortOptions::default()
+            && nulls == NullEquality::NullEqualsNothing
+        {
+            // Rows 4/5 equal the minimum; row 6 equals the maximum.
+            let expected = match op {
+                Operator::Lt => vec![4, 5],
+                Operator::Gt => vec![6],
+                Operator::NotEq | Operator::LtEq | Operator::GtEq => vec![4, 
5, 6],
+                _ => unreachable!(),
+            };
+            assert_eq!(outputs[1], expected, "{op:?}");
+        }
+    }
+    Ok(())
+}
+
+#[tokio::test]
+async fn guarded_or_preserves_anti_rows_and_requires_an_inner_witness() -> 
Result<()> {
+    let left = vec![
+        (Some(0), Some(5), Some(true)),
+        (Some(1), Some(2), Some(true)),
+        (Some(1), Some(3), Some(false)),
+        (Some(1), None, None),
+        (Some(2), Some(7), Some(true)),
+        (Some(3), None, Some(true)),
+    ];
+    let right = vec![
+        (Some(1), Some(2), Some(false)),
+        (Some(1), Some(4), Some(true)),
+        (Some(1), None, None),
+        (Some(2), Some(7), None),
+    ];
+    for kind in JOINS {
+        for outer_only_or in [false, true] {
+            let guarded = binary(
+                binary(
+                    column(2),
+                    Operator::And,
+                    binary(column(0), Operator::NotEq, column(1)),
+                ),
+                Operator::And,
+                column(3),
+            );
+            let expr = if outer_only_or {
+                binary(guarded, Operator::Or, 
Arc::new(IsNullExpr::new(column(0))))
+            } else {
+                binary(
+                    guarded,
+                    Operator::Or,
+                    binary(column(0), Operator::Lt, column(1)),
+                )
+            };
+            let expected = match (kind, outer_only_or) {
+                (JoinType::LeftSemi, false) => vec![1, 2],
+                (JoinType::LeftAnti, false) => vec![0, 3, 4, 5],
+                (JoinType::RightSemi, false) => vec![1],
+                (JoinType::RightAnti, false) => vec![0, 2, 3],
+                (JoinType::LeftSemi, true) => vec![1, 3],
+                (JoinType::LeftAnti, true) => vec![0, 2, 4, 5],
+                (JoinType::RightSemi, true) => vec![0, 1, 2],
+                (JoinType::RightAnti, true) => vec![3],
+                _ => unreachable!(),
+            };
+            for enabled in [false, true] {
+                let plan = join(
+                    input(&batch(&left)?, 1)?,
+                    input(&batch(&right)?, 2)?,
+                    kind,
+                    filter(Arc::clone(&expr), DataType::Int32),
+                    SortOptions::default(),
+                    NullEquality::NullEqualsNothing,
+                )?;
+                assert_eq!(
+                    ids(&plan, context(1, enabled)).await?,
+                    expected,
+                    "{kind:?} {outer_only_or} {enabled}"
+                );
+                assert_eq!(
+                    metric(&plan, "existence_summary_enabled"),
+                    usize::from(enabled)
+                );
+            }
+        }
+    }
+    Ok(())
+}
+
+#[tokio::test]
+async fn strings_dates_timestamps_and_decimals_match_generic_execution() -> 
Result<()> {
+    let rows = vec![(Some(1), None, Some(true)); 5];
+    let types: Vec<ArrayRef> = vec![
+        Arc::new(StringArray::from(vec![
+            None,
+            Some(""),
+            Some("é"),
+            Some("a"),
+            Some("é"),
+        ])),
+        Arc::new(LargeStringArray::from(vec![
+            None,
+            Some(""),
+            Some("é"),
+            Some("a"),
+            Some("é"),
+        ])),
+        Arc::new(StringViewArray::from(vec![
+            None,
+            Some(""),
+            Some("out-of-line-string-z"),
+            Some("out-of-line-string-a"),
+            Some("out-of-line-string-z"),
+        ])),
+        Arc::new(Date32Array::from(vec![
+            None,
+            Some(-1),
+            Some(0),
+            Some(1),
+            Some(1),
+        ])),
+        Arc::new(
+            TimestampMicrosecondArray::from(vec![
+                None,
+                Some(-100),
+                Some(0),
+                Some(100),
+                Some(100),
+            ])
+            .with_timezone("UTC"),
+        ),
+        Arc::new(
+            Decimal128Array::from(vec![None, Some(-123), Some(0), Some(456), 
Some(456)])
+                .with_precision_and_scale(20, 2)?,
+        ),
+    ];
+    for values in types {
+        let data_type = values.data_type().clone();
+        let batch = batch_values(&rows, values)?;
+        // NotEq uses scalar equality; ranges use array and scalar ordering.
+        for op in [Operator::NotEq, Operator::Lt] {
+            let mut outputs = vec![];
+            for enabled in [false, true] {
+                let plan = join(
+                    input(&batch, 2)?,
+                    input(&batch, 3)?,
+                    JoinType::LeftSemi,
+                    filter(binary(column(0), op, column(1)), 
data_type.clone()),
+                    SortOptions::default(),
+                    NullEquality::NullEqualsNothing,
+                )?;
+                outputs.push(ids(&plan, context(2, enabled)).await?);
+                assert_eq!(
+                    metric(&plan, "existence_summary_enabled"),
+                    usize::from(enabled),
+                    "{data_type:?} {op:?}"
+                );
+            }
+            assert_eq!(outputs[0], outputs[1], "{data_type:?} {op:?}");
+        }
+    }
+    Ok(())
+}
+
+#[tokio::test]
+async fn independent_cross_side_witnesses_use_generic_fallback() -> Result<()> 
{
+    let values = batch(&[
+        (Some(1), Some(0), Some(true)),
+        (Some(1), Some(5), None),
+        (Some(1), Some(10), Some(false)),
+    ])?;
+    // Min/max have separate witnesses for 5, but no row satisfies both 
clauses.
+    let expr = binary(
+        binary(column(0), Operator::Lt, column(1)),
+        Operator::And,
+        binary(column(0), Operator::Gt, column(1)),
+    );
+    for enabled in [false, true] {
+        let plan = join(
+            input(&values, 1)?,
+            input(&values, 2)?,
+            JoinType::LeftSemi,
+            filter(Arc::clone(&expr), DataType::Int32),
+            SortOptions::default(),
+            NullEquality::NullEqualsNothing,
+        )?;
+        assert!(ids(&plan, context(2, enabled)).await?.is_empty());
+        assert_eq!(metric(&plan, "existence_summary_enabled"), 0);
+        assert_eq!(
+            metric(&plan, "existence_summary_fallback"),
+            usize::from(enabled)
+        );
+    }
+    Ok(())
+}
+
+#[tokio::test]
+async fn normalized_nullable_strings_and_negated_equality_keep_sql_semantics()
+-> Result<()> {
+    let left = batch_values(
+        &[(Some(1), None, Some(true)); 4],
+        Arc::new(StringArray::from(vec![
+            None,
+            Some(""),
+            Some("a"),
+            Some("out-of-line-value"),
+        ])),
+    )?;
+    let right = batch_values(
+        &[(Some(1), None, Some(true)); 2],
+        Arc::new(StringArray::from(vec![None, Some("")])),
+    )?;
+    let normalize = |index| -> Result<Arc<dyn PhysicalExpr>> {
+        Ok(Arc::new(CaseExpr::try_new(
+            None,
+            vec![(
+                Arc::new(IsNullExpr::new(column(index))),
+                Arc::new(Literal::new(ScalarValue::Utf8(Some(String::new())))),
+            )],
+            Some(column(index)),
+        )?))
+    };
+    for kind in JOINS {
+        let expected = match kind {
+            JoinType::LeftSemi => vec![2, 3],
+            JoinType::LeftAnti => vec![0, 1],
+            JoinType::RightSemi => vec![0, 1],
+            JoinType::RightAnti => vec![],
+            _ => unreachable!(),
+        };
+        for enabled in [false, true] {
+            let expr = Arc::new(NotExpr::new(binary(
+                normalize(0)?,
+                Operator::Eq,
+                normalize(1)?,
+            )));
+            let plan = join(
+                input(&left, 1)?,
+                input(&right, 1)?,
+                kind,
+                filter(expr, DataType::Utf8),
+                SortOptions::default(),
+                NullEquality::NullEqualsNothing,
+            )?;
+            assert_eq!(
+                ids(&plan, context(1, enabled)).await?,
+                expected,
+                "{kind:?} {enabled}"
+            );
+            assert_eq!(
+                metric(&plan, "existence_summary_enabled"),
+                usize::from(enabled)
+            );
+        }
+    }
+    Ok(())
+}
+
+#[tokio::test]
+async fn empty_inputs_never_synthesize_a_witness() -> Result<()> {

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. Replaced the 
empty-input claim with a matching key whose seven inner operands are all NULL. 
It reaches reduction, produces null extrema, and preserves the anti rows with a 
false witness mask. The test asserts that summary storage was admitted 
alongside the buffer, remains reserved after releasing the buffer, and is 
released when the stream is dropped. An all-rejected guard belongs to ordinary 
evaluation now because guarded residuals are outside this version’s scope.



##########
datafusion/physical-plan/benches/sort_merge_join.rs:
##########
@@ -248,5 +251,167 @@ fn bench_smj(c: &mut Criterion) {
     group.finish();
 }
 
-criterion_group!(benches, bench_smj);
+/// Compare execution with summaries enabled and disabled in the same binary.
+/// Inputs are already sorted: SQL versions of these EXISTS/NOT EXISTS queries
+/// also measure sorting and depend on the optimizer's choice of join 
algorithm.
+/// These cases isolate the residual semi/anti join, including output 
collection.
+fn bench_existence_summary(c: &mut Criterion) {
+    let rt = Runtime::new().unwrap();
+    let mut group = c.benchmark_group("sort_merge_join_existence_summary");
+    group.sample_size(10);
+    group.warm_up_time(std::time::Duration::from_millis(500));
+    group.measurement_time(std::time::Duration::from_secs(2));
+
+    // Vary group count, rows per side, residual selectivity, and join type.
+    // The early-witness case controls for the generic join's short circuit;
+    // the small-group case measures the cost of repeatedly resetting 
summaries.
+    for (name, groups, probe_rows, inner_rows, distinct, op, kind) in [

Review Comment:
   Updated in 548a66268ab1928ed47b2cf778656b4386542815. Added the requested 
range-first-witness cases with one and many outer probes, plus a second-witness 
case, three-/four-/seven-row groups, 128-byte strings, guards, OR, and 
concurrent shared/separate FairSpillPool cases. The same harness runs on the 
exact PR base and this head, with output counts and zero-spill assertions. Also 
ran SMJ Q11–Q13/Q18, after verifying that their physical plans retain the 
eligible column comparison. The updated description records the validation 
scope and settings; full raw results remain local while the remaining cost 
questions are reviewed. I am leaving this thread open. 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 und
 er 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.



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