Skip to content

Commit d156e3e

Browse files
adriangbclaude
andcommitted
feat: add OptionalFilterPhysicalExpr to mark filters not needed for correctness
Add `OptionalFilterPhysicalExpr`, a transparent wrapper that marks a filter as optional: a consumer can skip it without changing the query result. A consumer can skip it only when the wrapper is a direct conjunct of the root AND chain of its predicate. In all other positions the wrapper is transparent, because `evaluate()` always evaluates the inner expression. `snapshot()` returns the inner expression, so pruning sees through it. Also add: - `split_optional` and `is_optional_filter` helpers in `physical_expr::utils` for consumers - `PhysicalOptionalFilterNode` proto message (field 29 in `PhysicalExprNode`) with self-encoding `try_to_proto`/`try_from_proto` No producer uses the wrapper yet, so there is no behavior change. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 22ce8bf commit d156e3e

10 files changed

Lines changed: 757 additions & 5 deletions

File tree

‎datafusion/physical-expr/src/expressions/mod.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ mod literal;
3333
mod negative;
3434
mod no_op;
3535
mod not;
36+
mod optional_filter;
3637
mod similar_to_pattern;
3738
mod try_cast;
3839
mod unknown_column;
@@ -59,6 +60,7 @@ pub use literal::{Literal, lit};
5960
pub use negative::{NegativeExpr, negative};
6061
pub use no_op::NoOp;
6162
pub use not::{NotExpr, not};
63+
pub use optional_filter::OptionalFilterPhysicalExpr;
6264
pub(crate) use similar_to_pattern::translate_scalar;
6365
pub use similar_to_pattern::{SqlSimilarToPattern, sql_similar_to_regex};
6466
pub use try_cast::{TryCastExpr, try_cast, try_cast_with_target_field};
Lines changed: 357 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,357 @@
1+
// Licensed to the Apache Software Foundation (ASF) under one
2+
// or more contributor license agreements. See the NOTICE file
3+
// distributed with this work for additional information
4+
// regarding copyright ownership. The ASF licenses this file
5+
// to you under the Apache License, Version 2.0 (the
6+
// "License"); you may not use this file except in compliance
7+
// with the License. You may obtain a copy of the License at
8+
//
9+
// http://www.apache.org/licenses/LICENSE-2.0
10+
//
11+
// Unless required by applicable law or agreed to in writing,
12+
// software distributed under the License is distributed on an
13+
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14+
// KIND, either express or implied. See the License for the
15+
// specific language governing permissions and limitations
16+
// under the License.
17+
18+
//! [`OptionalFilterPhysicalExpr`]: a marker for filters that are not needed
19+
//! for correctness. See the type documentation for the contract that
20+
//! producers, consumers and rewriters must obey.
21+
22+
use std::fmt;
23+
use std::hash::Hash;
24+
use std::sync::Arc;
25+
26+
use crate::PhysicalExpr;
27+
28+
use arrow::array::BooleanArray;
29+
use arrow::datatypes::{DataType, FieldRef, Schema};
30+
use arrow::record_batch::RecordBatch;
31+
use datafusion_common::{Result, assert_eq_or_internal_err};
32+
use datafusion_expr::ColumnarValue;
33+
use datafusion_expr::interval_arithmetic::Interval;
34+
use datafusion_expr::sort_properties::ExprProperties;
35+
36+
/// Marks the inner filter as *optional*: it is not needed for correctness.
37+
///
38+
/// Some filters are only performance hints. For example, a hash join can push
39+
/// a dynamic filter into the probe side scan, but the join itself still
40+
/// removes the rows that do not match. Such a filter is *optional*: a
41+
/// consumer can skip it (for example, when the filter does not remove enough
42+
/// rows to be worth its cost) and the query result stays the same.
43+
///
44+
/// # Contract
45+
///
46+
/// * **Skip only on the root AND chain.** A consumer can skip an optional
47+
/// filter only when the `Optional` node is a direct conjunct of the root
48+
/// `AND` chain of its predicate. For example, in `a AND Optional(b)` the
49+
/// consumer can skip `b`. Use [`split_optional`] to find these conjuncts.
50+
/// * **Transparent everywhere else.** [`PhysicalExpr::evaluate`] always
51+
/// evaluates the inner expression. Thus an `Optional` in a different
52+
/// position (for example under `NOT`, `IS NULL`, `CASE` or `OR`) can make a
53+
/// query slower, but it cannot make the result incorrect.
54+
/// * **Rewriters must not move nodes across the wrapper.** A rewrite must not
55+
/// move an expression into or out of an `Optional`. For example,
56+
/// `NOT(Optional(x))` must not become `Optional(NOT(x))`, because that
57+
/// would make a required filter optional.
58+
/// * **Pruning sees through the wrapper.** [`PhysicalExpr::snapshot`]
59+
/// returns the inner expression, so [`snapshot_physical_expr`] removes the
60+
/// wrapper. Thus statistics pruning uses an optional filter the same as a
61+
/// required filter.
62+
///
63+
/// [`split_optional`]: crate::utils::split_optional
64+
/// [`snapshot_physical_expr`]: datafusion_physical_expr_common::physical_expr::snapshot_physical_expr
65+
#[derive(Debug, Eq)]
66+
pub struct OptionalFilterPhysicalExpr {
67+
inner: Arc<dyn PhysicalExpr>,
68+
}
69+
70+
// Manually derive PartialEq and Hash to work around https://github.com/rust-lang/rust/issues/78808
71+
impl PartialEq for OptionalFilterPhysicalExpr {
72+
fn eq(&self, other: &Self) -> bool {
73+
self.inner.eq(&other.inner)
74+
}
75+
}
76+
77+
impl Hash for OptionalFilterPhysicalExpr {
78+
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
79+
self.inner.hash(state);
80+
}
81+
}
82+
83+
impl OptionalFilterPhysicalExpr {
84+
/// Create a new optional filter that wraps `inner`.
85+
pub fn new(inner: Arc<dyn PhysicalExpr>) -> Self {
86+
Self { inner }
87+
}
88+
89+
/// Get the wrapped filter expression.
90+
pub fn inner(&self) -> &Arc<dyn PhysicalExpr> {
91+
&self.inner
92+
}
93+
}
94+
95+
impl fmt::Display for OptionalFilterPhysicalExpr {
96+
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
97+
write!(f, "Optional({})", self.inner)
98+
}
99+
}
100+
101+
impl PhysicalExpr for OptionalFilterPhysicalExpr {
102+
fn data_type(&self, input_schema: &Schema) -> Result<DataType> {
103+
self.inner.data_type(input_schema)
104+
}
105+
106+
fn nullable(&self, input_schema: &Schema) -> Result<bool> {
107+
self.inner.nullable(input_schema)
108+
}
109+
110+
fn evaluate(&self, batch: &RecordBatch) -> Result<ColumnarValue> {
111+
self.inner.evaluate(batch)
112+
}
113+
114+
fn return_field(&self, input_schema: &Schema) -> Result<FieldRef> {
115+
self.inner.return_field(input_schema)
116+
}
117+
118+
fn evaluate_selection(
119+
&self,
120+
batch: &RecordBatch,
121+
selection: &BooleanArray,
122+
) -> Result<ColumnarValue> {
123+
self.inner.evaluate_selection(batch, selection)
124+
}
125+
126+
fn children(&self) -> Vec<&Arc<dyn PhysicalExpr>> {
127+
vec![&self.inner]
128+
}
129+
130+
fn with_new_children(
131+
self: Arc<Self>,
132+
children: Vec<Arc<dyn PhysicalExpr>>,
133+
) -> Result<Arc<dyn PhysicalExpr>> {
134+
assert_eq_or_internal_err!(
135+
children.len(),
136+
1,
137+
"OptionalFilterPhysicalExpr: expected 1 child"
138+
);
139+
Ok(Arc::new(Self::new(Arc::clone(&children[0]))))
140+
}
141+
142+
// The wrapper is the identity function, so the bounds and properties of
143+
// the child are also the bounds and properties of the wrapper.
144+
fn evaluate_bounds(&self, children: &[&Interval]) -> Result<Interval> {
145+
Ok(children[0].clone())
146+
}
147+
148+
fn propagate_constraints(
149+
&self,
150+
interval: &Interval,
151+
children: &[&Interval],
152+
) -> Result<Option<Vec<Interval>>> {
153+
Ok(children[0].intersect(interval)?.map(|result| vec![result]))
154+
}
155+
156+
fn get_properties(&self, children: &[ExprProperties]) -> Result<ExprProperties> {
157+
Ok(children[0].clone())
158+
}
159+
160+
fn fmt_sql(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
161+
self.inner.fmt_sql(f)
162+
}
163+
164+
/// Returns the inner expression, so that snapshot consumers (for example
165+
/// pruning) see through the wrapper.
166+
///
167+
/// [`snapshot_physical_expr`] transforms the tree bottom up, so the inner
168+
/// expression is already a snapshot when this method is called. For
169+
/// example, `Optional(DynamicFilter)` becomes the current expression of
170+
/// the dynamic filter.
171+
///
172+
/// [`snapshot_physical_expr`]: datafusion_physical_expr_common::physical_expr::snapshot_physical_expr
173+
fn snapshot(&self) -> Result<Option<Arc<dyn PhysicalExpr>>> {
174+
Ok(Some(Arc::clone(&self.inner)))
175+
}
176+
177+
fn snapshot_generation(&self) -> u64 {
178+
// The wrapper is not dynamic. `snapshot_generation(expr)` walks the
179+
// tree and adds the generation of the inner expression.
180+
0
181+
}
182+
183+
#[cfg(feature = "proto")]
184+
fn try_to_proto(
185+
&self,
186+
ctx: &datafusion_physical_expr_common::physical_expr::proto_encode::PhysicalExprEncodeCtx<'_>,
187+
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
188+
use datafusion_proto_models::protobuf;
189+
190+
Ok(Some(protobuf::PhysicalExprNode {
191+
expr_id: None,
192+
expr_type: Some(protobuf::physical_expr_node::ExprType::OptionalFilter(
193+
Box::new(protobuf::PhysicalOptionalFilterNode {
194+
inner: Some(Box::new(ctx.encode_child(&self.inner)?)),
195+
}),
196+
)),
197+
}))
198+
}
199+
}
200+
201+
#[cfg(feature = "proto")]
202+
impl OptionalFilterPhysicalExpr {
203+
/// Reconstruct an [`OptionalFilterPhysicalExpr`] from its protobuf
204+
/// representation.
205+
pub fn try_from_proto(
206+
node: &datafusion_proto_models::protobuf::PhysicalExprNode,
207+
ctx: &datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx<'_>,
208+
) -> Result<Arc<dyn PhysicalExpr>> {
209+
use datafusion_physical_expr_common::expect_expr_variant;
210+
use datafusion_proto_models::protobuf;
211+
212+
let optional = expect_expr_variant!(
213+
node,
214+
protobuf::physical_expr_node::ExprType::OptionalFilter,
215+
"OptionalFilter",
216+
);
217+
let inner = ctx.decode_required_expression(
218+
optional.inner.as_deref(),
219+
"OptionalFilter",
220+
"inner",
221+
)?;
222+
223+
Ok(Arc::new(Self::new(inner)))
224+
}
225+
}
226+
227+
#[cfg(test)]
228+
mod tests {
229+
use super::*;
230+
use crate::expressions::{BinaryExpr, DynamicFilterPhysicalExpr, col, lit, not};
231+
232+
use arrow::array::{ArrayRef, Int32Array};
233+
use arrow::datatypes::Field;
234+
use datafusion_common::cast::as_boolean_array;
235+
use datafusion_expr::Operator;
236+
use datafusion_physical_expr_common::physical_expr::{
237+
fmt_sql, snapshot_generation, snapshot_physical_expr,
238+
};
239+
240+
fn schema() -> Arc<Schema> {
241+
Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, true)]))
242+
}
243+
244+
/// `a > 2`
245+
fn a_gt_2(schema: &Schema) -> Arc<dyn PhysicalExpr> {
246+
Arc::new(BinaryExpr::new(
247+
col("a", schema).unwrap(),
248+
Operator::Gt,
249+
lit(2i32),
250+
))
251+
}
252+
253+
fn optional(inner: Arc<dyn PhysicalExpr>) -> Arc<dyn PhysicalExpr> {
254+
Arc::new(OptionalFilterPhysicalExpr::new(inner))
255+
}
256+
257+
#[test]
258+
fn evaluate_equals_inner() -> Result<()> {
259+
let schema = schema();
260+
let values: ArrayRef = Arc::new(Int32Array::from(vec![Some(1), None, Some(3)]));
261+
let batch = RecordBatch::try_new(Arc::clone(&schema), vec![values])?;
262+
263+
let inner = a_gt_2(&schema);
264+
let wrapped = optional(Arc::clone(&inner));
265+
266+
let expected = inner.evaluate(&batch)?.into_array(batch.num_rows())?;
267+
let actual = wrapped.evaluate(&batch)?.into_array(batch.num_rows())?;
268+
assert_eq!(as_boolean_array(&actual)?, as_boolean_array(&expected)?);
269+
270+
let selection = BooleanArray::from(vec![true, false, true]);
271+
let expected = inner
272+
.evaluate_selection(&batch, &selection)?
273+
.into_array(batch.num_rows())?;
274+
let actual = wrapped
275+
.evaluate_selection(&batch, &selection)?
276+
.into_array(batch.num_rows())?;
277+
assert_eq!(as_boolean_array(&actual)?, as_boolean_array(&expected)?);
278+
279+
assert_eq!(wrapped.data_type(&schema)?, DataType::Boolean);
280+
assert!(wrapped.nullable(&schema)?);
281+
Ok(())
282+
}
283+
284+
#[test]
285+
fn display_and_fmt_sql() {
286+
let schema = schema();
287+
let wrapped = optional(a_gt_2(&schema));
288+
assert_eq!(wrapped.to_string(), "Optional(a@0 > 2)");
289+
assert_eq!(fmt_sql(wrapped.as_ref()).to_string(), "a > 2");
290+
291+
let negated = not(Arc::clone(&wrapped)).unwrap();
292+
assert_eq!(negated.to_string(), "NOT Optional(a@0 > 2)");
293+
}
294+
295+
#[test]
296+
fn children_and_with_new_children() -> Result<()> {
297+
let schema = schema();
298+
let wrapped = optional(a_gt_2(&schema));
299+
assert_eq!(wrapped.children().len(), 1);
300+
301+
let new_inner = lit(true);
302+
let rewrapped = Arc::clone(&wrapped).with_new_children(vec![new_inner])?;
303+
let rewrapped = rewrapped
304+
.downcast_ref::<OptionalFilterPhysicalExpr>()
305+
.expect("wrapper is kept");
306+
assert_eq!(rewrapped.inner().to_string(), "true");
307+
308+
assert!(wrapped.with_new_children(vec![]).is_err());
309+
Ok(())
310+
}
311+
312+
#[test]
313+
fn eq_and_hash_use_inner() {
314+
use std::collections::HashSet;
315+
316+
let schema = schema();
317+
let a = optional(a_gt_2(&schema));
318+
let b = optional(a_gt_2(&schema));
319+
let c = optional(lit(true));
320+
assert_eq!(&a, &b);
321+
assert_ne!(&a, &c);
322+
// The wrapper is not equal to the inner expression.
323+
assert_ne!(&a, &a_gt_2(&schema));
324+
325+
let set: HashSet<_> = [a, b, c].into_iter().collect();
326+
assert_eq!(set.len(), 2);
327+
}
328+
329+
#[test]
330+
fn snapshot_sees_through_dynamic_filter() -> Result<()> {
331+
let schema = schema();
332+
let dynamic = Arc::new(DynamicFilterPhysicalExpr::new(
333+
vec![col("a", &schema)?],
334+
lit(true),
335+
));
336+
let wrapped = optional(Arc::clone(&dynamic) as Arc<dyn PhysicalExpr>);
337+
338+
assert_eq!(
339+
snapshot_physical_expr(Arc::clone(&wrapped))?.to_string(),
340+
"true"
341+
);
342+
343+
let generation = snapshot_generation(&wrapped);
344+
dynamic.update(a_gt_2(&schema))?;
345+
assert_ne!(snapshot_generation(&wrapped), generation);
346+
assert_eq!(
347+
snapshot_physical_expr(Arc::clone(&wrapped))?.to_string(),
348+
"a@0 > 2"
349+
);
350+
351+
// A static inner expression is also unwrapped.
352+
let wrapped = optional(a_gt_2(&schema));
353+
assert_eq!(wrapped.snapshot_generation(), 0);
354+
assert_eq!(snapshot_physical_expr(wrapped)?.to_string(), "a@0 > 2");
355+
Ok(())
356+
}
357+
}

0 commit comments

Comments
 (0)