viirya commented on code in PR #24526:
URL: https://github.com/apache/datafusion/pull/24526#discussion_r3864684260


##########
datafusion/pruning/src/string_in_list.rs:
##########
@@ -0,0 +1,225 @@
+// 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.
+
+use std::fmt::{self, Display, Formatter};
+use std::hash::{Hash, Hasher};
+use std::sync::Arc;
+
+use arrow::array::{Array, AsArray, BooleanArray};
+use arrow::compute::cast;
+use arrow::datatypes::{DataType, Schema};
+use arrow::record_batch::RecordBatch;
+use datafusion_common::{Result, assert_eq_or_internal_err};
+use datafusion_physical_expr::{PhysicalExpr, PhysicalExprRef};
+use datafusion_physical_plan::ColumnarValue;
+
+/// Tests whether a sorted string domain intersects an inclusive statistics 
interval.
+/// This expression is used only for pruning; the original IN remains the row 
filter.
+#[derive(Debug, Eq)]
+pub(crate) struct StringInListPruningExpr {
+    min: PhysicalExprRef,
+    max: PhysicalExprRef,
+    values: Arc<[String]>,
+}
+
+impl StringInListPruningExpr {
+    pub(crate) fn new(
+        min: PhysicalExprRef,
+        max: PhysicalExprRef,
+        mut values: Vec<String>,
+    ) -> Self {
+        values.sort_unstable();
+        values.dedup();
+        Self {
+            min,
+            max,
+            values: values.into(),
+        }
+    }
+}
+
+impl PartialEq for StringInListPruningExpr {
+    fn eq(&self, other: &Self) -> bool {
+        self.min.eq(&other.min) && self.max.eq(&other.max) && self.values == 
other.values
+    }
+}
+
+impl Hash for StringInListPruningExpr {
+    fn hash<H: Hasher>(&self, state: &mut H) {
+        self.min.hash(state);
+        self.max.hash(state);
+        self.values.hash(state);
+    }
+}
+
+impl Display for StringInListPruningExpr {
+    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
+        write!(
+            f,
+            "IN_SET_INTERSECTS({}, {}, {} values)",
+            self.min,
+            self.max,
+            self.values.len()
+        )
+    }
+}
+
+fn has_oversized_string_buffer(array: &dyn Array, limit: usize) -> bool {
+    match array.data_type() {
+        DataType::Utf8 => array.as_string::<i32>().values().len() >= limit,
+        DataType::LargeUtf8 => array.as_string::<i64>().values().len() >= 
limit,
+        DataType::Dictionary(_, _) => has_oversized_string_buffer(
+            array.as_any_dictionary().values().as_ref(),
+            limit,
+        ),
+        _ => false,
+    }
+}
+
+impl PhysicalExpr for StringInListPruningExpr {
+    fn data_type(&self, _input_schema: &Schema) -> Result<DataType> {
+        Ok(DataType::Boolean)
+    }
+
+    fn nullable(&self, _input_schema: &Schema) -> Result<bool> {
+        Ok(true)
+    }
+
+    fn evaluate(&self, batch: &RecordBatch) -> Result<ColumnarValue> {
+        // Normalize Utf8, LargeUtf8, Utf8View, and dictionary-encoded 
statistics.
+        let min = self.min.evaluate(batch)?.into_array(batch.num_rows())?;
+        let max = self.max.evaluate(batch)?.into_array(batch.num_rows())?;
+        // A short string slice can retain a buffer too large for Utf8View's
+        // u32 offsets. Avoid a panic in the cast and keep pruning 
conservative.
+        if has_oversized_string_buffer(min.as_ref(), u32::MAX as usize)
+            || has_oversized_string_buffer(max.as_ref(), u32::MAX as usize)
+        {
+            return Ok(ColumnarValue::Array(Arc::new(BooleanArray::new_null(
+                batch.num_rows(),
+            ))));
+        }
+        // Dictionary values can be NULL behind valid keys. Preserve their
+        // validity even if the view cast only carries the key nulls.
+        let min_nulls = min.logical_nulls();
+        let max_nulls = max.logical_nulls();
+        let min = cast(&min, &DataType::Utf8View)?;
+        let max = cast(&max, &DataType::Utf8View)?;
+        let min = min.as_string_view();
+        let max = max.as_string_view();
+        let matches: BooleanArray = (0..batch.num_rows())
+            .map(|i| {
+                let min = (min.is_valid(i)
+                    && min_nulls.as_ref().is_none_or(|nulls| 
nulls.is_valid(i)))
+                .then(|| min.value(i).as_bytes());
+                let max = (max.is_valid(i)
+                    && max_nulls.as_ref().is_none_or(|nulls| 
nulls.is_valid(i)))
+                .then(|| max.value(i).as_bytes());
+                match (min, max) {
+                    (Some(min), Some(max)) => {
+                        if min > max {
+                            return None;
+                        }
+                        // Rust string ordering and these byte comparisons 
both use
+                        // unsigned lexicographic UTF-8 order. Statistics 
providers

Review Comment:
   This comment states a contract on statistics providers generally, but the 
enforcement it points at is Parquet-specific. I checked the other providers 
that can reach this expression — eligibility keys only on the schema type being 
a string, so it carries no knowledge of who supplied the bounds:
   
   - `PartitionPruningStatistics` — bounds are partition values DataFusion 
produced itself, so they are in Arrow order by construction.
   - `PrunableStatistics` (built in `file_pruner.rs` from a file's 
`ColumnStatistics`) — bounds come from whatever the `TableProvider` reported, 
with no ordering gate on that path.
   
   The exposure isn't new; the per-value `min <= v AND v <= max` path has the 
same dependency. But since this is the first place the requirement is written 
down, the comment should say which providers are known to satisfy it and which 
are taken on trust — as written it reads as if the invariant is enforced 
everywhere. (If you think it belongs closer to the source, 
`PruningStatistics::min_values` might be the better home for it.)



##########
datafusion/pruning/src/pruning_predicate.rs:
##########
@@ -1582,6 +1629,16 @@ fn build_predicate_expression(
         }
     }
     if let Some(in_list) = expr.downcast_ref::<phys_expr::InListExpr>() {
+        // Preserve the existing per-value representation for lists within
+        // the default limit. Use the compact form only when callers raise
+        // the cap; MAX_IN_LIST_SIZE is not a measured performance crossover.
+        if in_list.list().len() > MAX_IN_LIST_SIZE

Review Comment:
   This is the one remaining item from the earlier thread — you agreed the 
scope/default-behaviour rationale should sit next to this condition, and it 
doesn't appear to be in `df65b7cc` yet. A sentence saying the lower bound is a 
scope/compatibility choice rather than a measured threshold would stop the next 
reader assuming 20 is tuned.



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