From 8d3ac7b34d99b87a16e8731f6c7eee9a29c6bc40 Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Fri, 6 Feb 2026 11:21:01 +0100 Subject: [PATCH 01/15] wip --- .../src/plans/aexpr/function_expr/mod.rs | 9 ++++ crates/polars-plan/src/plans/optimizer/mod.rs | 2 +- .../optimizer/predicate_pushdown/dynamic.rs | 45 +++++++++++++++++++ .../plans/optimizer/predicate_pushdown/mod.rs | 2 + 4 files changed, 57 insertions(+), 1 deletion(-) create mode 100644 crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs diff --git a/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs b/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs index cdec10498a46..bbab33ed4ea2 100644 --- a/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs +++ b/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs @@ -87,6 +87,7 @@ pub use self::struct_::IRStructFunction; #[cfg(feature = "trigonometry")] pub use self::trigonometry::IRTrigonometricFunction; use super::*; +use crate::plans::optimizer::DynamicPred; #[cfg_attr(feature = "ir_serde", derive(serde::Serialize, serde::Deserialize))] #[derive(Clone, PartialEq, Debug)] @@ -387,6 +388,9 @@ pub enum IRFunctionExpr { RowEncode(Vec, RowEncodingVariant), #[cfg(feature = "dtype-struct")] RowDecode(Vec, RowEncodingVariant), + DynamicExpr { + pred: DynamicPred, + }, } impl Hash for IRFunctionExpr { @@ -690,6 +694,9 @@ impl Hash for IRFunctionExpr { fs.hash(state); variants.hash(state); }, + DynamicExpr { pred } => { + pred.id().hash(state); + }, } } } @@ -904,6 +911,7 @@ impl Display for IRFunctionExpr { RowEncode(..) => "row_encode", #[cfg(feature = "dtype-struct")] RowDecode(..) => "row_decode", + DynamicExpr { pred } => "dynamic_predicate", }; write!(f, "{s}") } @@ -1233,6 +1241,7 @@ impl IRFunctionExpr { F::RowEncode(..) => FunctionOptions::elementwise(), #[cfg(feature = "dtype-struct")] F::RowDecode(..) => FunctionOptions::elementwise(), + F::DynamicExpr { .. } => FunctionOptions::elementwise(), } } } diff --git a/crates/polars-plan/src/plans/optimizer/mod.rs b/crates/polars-plan/src/plans/optimizer/mod.rs index 1273ab8e1342..21a0e52e272d 100644 --- a/crates/polars-plan/src/plans/optimizer/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/mod.rs @@ -34,7 +34,7 @@ pub use cse::NaiveExprMerger; use delay_rechunk::DelayRechunk; pub use expand_datasets::ExpandedDataset; use polars_core::config::verbose; -pub use predicate_pushdown::PredicatePushDown; +pub use predicate_pushdown::{DynamicPred, PredicatePushDown}; pub use projection_pushdown::ProjectionPushDown; pub use simplify_expr::{SimplifyBooleanRule, SimplifyExprRule}; use slice_pushdown_lp::SlicePushDown; diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs new file mode 100644 index 000000000000..881a13447d66 --- /dev/null +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -0,0 +1,45 @@ +use std::fmt::{Debug, Formatter}; +use std::sync::RwLock; +use std::sync::atomic::AtomicBool; + +use polars_utils::unique_id::UniqueId; +#[cfg(feature = "serde")] +use serde::{Deserialize, Serialize}; + +use super::*; + +#[cfg_attr(feature = "ir_serde", derive(Serialize, Deserialize))] +struct Inner { + pred: RwLock>, + is_set: AtomicBool, + id: UniqueId, +} + +#[derive(Clone)] +#[cfg_attr(feature = "ir_serde", derive(Serialize, Deserialize))] +pub struct DynamicPred { + inner: Arc, +} + +impl Debug for DynamicPred { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "dynamic_pred: {:}", self.id()) + } +} + +impl PartialEq for DynamicPred { + fn eq(&self, other: &Self) -> bool { + self.id() == other.id() + } +} + +impl DynamicPred { + pub fn id(&self) -> &UniqueId { + &self.inner.id + } + + fn set(&self, node: Node) { + let mut guard = self.inner.pred.write().unwrap(); + *guard = Some(node); + } +} diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs index 5cdf9ed37358..e4c0d5a9d6ae 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs @@ -1,8 +1,10 @@ +mod dynamic; mod group_by; mod join; mod keys; mod utils; +pub use dynamic::DynamicPred; use polars_core::datatypes::PlHashMap; use polars_core::prelude::*; use polars_utils::idx_vec::UnitVec; From 1b08cd7234d49e366050689710d44ec27330f7d2 Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Fri, 6 Feb 2026 14:44:07 +0100 Subject: [PATCH 02/15] dispatch --- crates/polars-expr/src/dispatch/misc.rs | 6 ++++- crates/polars-expr/src/dispatch/mod.rs | 3 +++ .../src/plans/aexpr/function_expr/schema.rs | 1 + .../src/plans/conversion/ir_to_dsl.rs | 8 ++++++ .../optimizer/predicate_pushdown/dynamic.rs | 25 ++++++++++++++++--- 5 files changed, 38 insertions(+), 5 deletions(-) diff --git a/crates/polars-expr/src/dispatch/misc.rs b/crates/polars-expr/src/dispatch/misc.rs index 6f224d57fd24..03527ea36455 100644 --- a/crates/polars-expr/src/dispatch/misc.rs +++ b/crates/polars-expr/src/dispatch/misc.rs @@ -16,7 +16,7 @@ use polars_plan::dsl::ReshapeDimension; use polars_plan::plans::FusedOperator; #[cfg(feature = "cov")] use polars_plan::plans::IRCorrelationMethod; -use polars_plan::plans::RowEncodingVariant; +use polars_plan::plans::{DynamicPred, RowEncodingVariant}; use polars_row::RowEncodingOptions; use polars_utils::IdxSize; use polars_utils::pl_str::PlSmallStr; @@ -1054,3 +1054,7 @@ pub fn repeat(args: &[Column]) -> PolarsResult { Ok(c.new_from_index(0, n)) } + +pub fn dynamic_expr(columns: &[Column], pred: &DynamicPred) -> PolarsResult { + pred.evaluate(columns) +} diff --git a/crates/polars-expr/src/dispatch/mod.rs b/crates/polars-expr/src/dispatch/mod.rs index 1d15610a5b16..90a4dc68f84d 100644 --- a/crates/polars-expr/src/dispatch/mod.rs +++ b/crates/polars-expr/src/dispatch/mod.rs @@ -527,6 +527,9 @@ pub fn function_expr_to_udf(func: IRFunctionExpr) -> SpecialEq { map_as_slice!(misc::row_decode, fs.clone(), variants.clone()) }, + F::DynamicExpr { pred } => { + map_as_slice!(misc::dynamic_expr, &pred) + }, } } diff --git a/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs b/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs index 8e70a24476f4..49a91adfa986 100644 --- a/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs +++ b/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs @@ -451,6 +451,7 @@ impl IRFunctionExpr { }), #[cfg(feature = "dtype-struct")] RowDecode(fields, _) => mapper.with_dtype(DataType::Struct(fields.to_vec())), + DynamicExpr { .. } => mapper.with_dtype(DataType::Boolean), } } diff --git a/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs b/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs index 8c786a132dae..9121db6ce863 100644 --- a/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs +++ b/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs @@ -1,3 +1,5 @@ +use polars_utils::format_pl_smallstr; + use super::*; /// converts a node from the AExpr arena to Expr @@ -1175,6 +1177,12 @@ pub fn ir_function_to_dsl(input: Vec, function: IRFunctionExpr) -> Expr { fs.into_iter().map(|f| (f.name, f.dtype.into())).collect(), v, ), + IF::DynamicExpr { pred } => { + return Expr::Display { + inputs: input, + fmt_str: Box::new(format_pl_smallstr!("{pred:?}")), + }; + }, }; Expr::Function { input, function } diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 881a13447d66..30335e730d17 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -2,6 +2,7 @@ use std::fmt::{Debug, Formatter}; use std::sync::RwLock; use std::sync::atomic::AtomicBool; +use polars_core::frame::column::ScalarColumn; use polars_utils::unique_id::UniqueId; #[cfg(feature = "serde")] use serde::{Deserialize, Serialize}; @@ -10,7 +11,9 @@ use super::*; #[cfg_attr(feature = "ir_serde", derive(Serialize, Deserialize))] struct Inner { - pred: RwLock>, + #[cfg_attr(feature = "ir_serde", serde(skip))] + pred: RwLock PolarsResult + Send + Sync>>>, + #[cfg_attr(feature = "ir_serde", serde(skip))] is_set: AtomicBool, id: UniqueId, } @@ -38,8 +41,22 @@ impl DynamicPred { &self.inner.id } - fn set(&self, node: Node) { - let mut guard = self.inner.pred.write().unwrap(); - *guard = Some(node); + pub fn evaluate(&self, columns: &[Column]) -> PolarsResult { + polars_ensure!(columns.len() > 1, ComputeError: "expected at least 1 argument"); + + let h = columns[0].len(); + + if self.inner.is_set.load(std::sync::atomic::Ordering::Relaxed) { + let guard = self.inner.pred.read().unwrap(); + let dyn_func = guard.as_ref().unwrap(); + dyn_func(columns) + } else { + let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(true)); + Ok(Column::Scalar(ScalarColumn::new( + columns[0].name().clone(), + s, + h, + ))) + } } } From 79e4d3c2f018370090b83667ac83ea317771b857 Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Fri, 6 Feb 2026 15:40:34 +0100 Subject: [PATCH 03/15] set --- .../src/plans/optimizer/predicate_pushdown/dynamic.rs | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 30335e730d17..6ecc04f36725 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -41,6 +41,16 @@ impl DynamicPred { &self.inner.id } + pub fn set(&self, pred: Box PolarsResult + Send + Sync>) { + { + let mut guard = self.inner.pred.write().unwrap(); + *guard = Some(pred); + } + self.inner + .is_set + .store(true, std::sync::atomic::Ordering::Release); + } + pub fn evaluate(&self, columns: &[Column]) -> PolarsResult { polars_ensure!(columns.len() > 1, ComputeError: "expected at least 1 argument"); From b5facf987cb599ed8d0aac7c09c84f272e95d0db Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Sat, 7 Feb 2026 09:59:26 +0100 Subject: [PATCH 04/15] create dyn pred --- .../optimizer/predicate_pushdown/dynamic.rs | 23 ++++++++++++++++++ .../plans/optimizer/predicate_pushdown/mod.rs | 24 ++++++++++++++++++- 2 files changed, 46 insertions(+), 1 deletion(-) diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 6ecc04f36725..3524becb3d71 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -37,6 +37,16 @@ impl PartialEq for DynamicPred { } impl DynamicPred { + fn new() -> Self { + Self { + inner: Arc::new(Inner { + pred: Default::default(), + is_set: Default::default(), + id: UniqueId::new(), + }), + } + } + pub fn id(&self) -> &UniqueId { &self.inner.id } @@ -70,3 +80,16 @@ impl DynamicPred { } } } + +pub fn new_dynamic_pred(node: Node, arena: &mut Arena) -> (Node, DynamicPred) { + let pred = DynamicPred::new(); + let function = IRFunctionExpr::DynamicExpr { pred: pred.clone() }; + let options = function.function_options(); + let aexpr = AExpr::Function { + input: vec![ExprIR::from_node(node, arena)], + function, + options, + }; + + (arena.add(aexpr), pred) +} diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs index e4c0d5a9d6ae..f02beea5e071 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs @@ -12,6 +12,7 @@ use recursive::recursive; use utils::*; use super::*; +use crate::plans::optimizer::predicate_pushdown::dynamic::new_dynamic_pred; use crate::prelude::optimizer::predicate_pushdown::group_by::process_group_by; use crate::prelude::optimizer::predicate_pushdown::join::process_join; use crate::utils::{check_input_node, has_aexpr}; @@ -574,7 +575,28 @@ impl PredicatePushDown { lp @ Union { .. } => { self.pushdown_and_continue(lp, acc_predicates, lp_arena, expr_arena, false) }, - lp @ Sort { .. } => { + Sort { + input, + by_column, + slice, + sort_options, + } => { + if let Some((offset, len)) = slice + && by_column.len() == 1 + { + let n = by_column[0].node(); + if let AExpr::Column(_) = expr_arena.get(n) { + let (node, pred) = new_dynamic_pred(n, expr_arena); + // todo + } + } + + let lp = Sort { + input, + by_column, + slice, + sort_options, + }; self.pushdown_and_continue(lp, acc_predicates, lp_arena, expr_arena, true) }, lp @ Sink { .. } | lp @ SinkMultiple { .. } => { From 37aa9d30099d050838b6f5021d0d897a5800fa80 Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Sat, 7 Feb 2026 10:22:07 +0100 Subject: [PATCH 05/15] call stack --- crates/polars-mem-engine/src/planner/lp.rs | 2 +- crates/polars-plan/src/plans/builder_ir.rs | 2 +- crates/polars-plan/src/plans/conversion/dsl_to_ir/mod.rs | 2 +- crates/polars-plan/src/plans/ir/format.rs | 8 ++++++-- crates/polars-plan/src/plans/ir/mod.rs | 2 +- .../src/plans/optimizer/predicate_pushdown/dynamic.rs | 7 +++++++ .../src/plans/optimizer/predicate_pushdown/mod.rs | 2 +- .../polars-plan/src/plans/optimizer/slice_pushdown_lp.rs | 2 +- crates/polars-stream/src/physical_plan/fmt.rs | 1 + crates/polars-stream/src/physical_plan/lower_expr.rs | 1 + crates/polars-stream/src/physical_plan/lower_ir.rs | 7 ++++--- crates/polars-stream/src/physical_plan/mod.rs | 3 ++- crates/polars-stream/src/physical_plan/to_graph.rs | 3 ++- 13 files changed, 29 insertions(+), 13 deletions(-) diff --git a/crates/polars-mem-engine/src/planner/lp.rs b/crates/polars-mem-engine/src/planner/lp.rs index cd46da1ac82d..a63a864b7900 100644 --- a/crates/polars-mem-engine/src/planner/lp.rs +++ b/crates/polars-mem-engine/src/planner/lp.rs @@ -451,7 +451,7 @@ fn create_physical_plan_impl( Ok(Box::new(executors::SortExec { input, by_column, - slice, + slice: slice.map(|t| (t.0, t.1)), sort_options, })) }, diff --git a/crates/polars-plan/src/plans/builder_ir.rs b/crates/polars-plan/src/plans/builder_ir.rs index bac26169c6c5..aab7eeedfb71 100644 --- a/crates/polars-plan/src/plans/builder_ir.rs +++ b/crates/polars-plan/src/plans/builder_ir.rs @@ -184,7 +184,7 @@ impl<'a> IRBuilder<'a> { let ir = IR::Sort { input: self.root, by_column, - slice, + slice: slice.map(|t| (t.0, t.1, None)), sort_options, }; let node = self.lp_arena.add(ir); diff --git a/crates/polars-plan/src/plans/conversion/dsl_to_ir/mod.rs b/crates/polars-plan/src/plans/conversion/dsl_to_ir/mod.rs index b2892b6de9bc..edbffa4c2783 100644 --- a/crates/polars-plan/src/plans/conversion/dsl_to_ir/mod.rs +++ b/crates/polars-plan/src/plans/conversion/dsl_to_ir/mod.rs @@ -506,7 +506,7 @@ pub fn to_alp_impl(lp: DslPlan, ctxt: &mut DslConversionContext) -> PolarsResult let lp = IR::Sort { input, by_column, - slice, + slice: slice.map(|t| (t.0, t.1, None)), sort_options, }; diff --git a/crates/polars-plan/src/plans/ir/format.rs b/crates/polars-plan/src/plans/ir/format.rs index 4079c7158558..89adb3a7c6c7 100644 --- a/crates/polars-plan/src/plans/ir/format.rs +++ b/crates/polars-plan/src/plans/ir/format.rs @@ -851,8 +851,12 @@ pub fn write_ir_non_recursive( f.write_char('[')?; let mut comma = false; - if let Some((o, l)) = slice { - write!(f, "slice: ({o}, {l})")?; + if let Some((o, l, dyn_pred)) = slice { + if let Some(dyn_pred) = &dyn_pred { + write!(f, "slice: ({o}, {l})")?; + } else { + write!(f, "slice: ({o}, {l}, {dyn_pred:?})")?; + } comma = true; } if sort_options.maintain_order { diff --git a/crates/polars-plan/src/plans/ir/mod.rs b/crates/polars-plan/src/plans/ir/mod.rs index 9e5db7982a71..c8dbc98216b7 100644 --- a/crates/polars-plan/src/plans/ir/mod.rs +++ b/crates/polars-plan/src/plans/ir/mod.rs @@ -89,7 +89,7 @@ pub enum IR { Sort { input: Node, by_column: Vec, - slice: Option<(i64, usize)>, + slice: Option<(i64, usize, Option)>, sort_options: SortMultipleOptions, }, Cache { diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 3524becb3d71..c24137a7a77b 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -1,4 +1,5 @@ use std::fmt::{Debug, Formatter}; +use std::hash::Hash; use std::sync::RwLock; use std::sync::atomic::AtomicBool; @@ -36,6 +37,12 @@ impl PartialEq for DynamicPred { } } +impl Hash for DynamicPred { + fn hash(&self, state: &mut H) { + self.inner.id.hash(state); + } +} + impl DynamicPred { fn new() -> Self { Self { diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs index f02beea5e071..f46eb7f55610 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs @@ -581,7 +581,7 @@ impl PredicatePushDown { slice, sort_options, } => { - if let Some((offset, len)) = slice + if let Some((offset, len, None)) = slice && by_column.len() == 1 { let n = by_column[0].node(); diff --git a/crates/polars-plan/src/plans/optimizer/slice_pushdown_lp.rs b/crates/polars-plan/src/plans/optimizer/slice_pushdown_lp.rs index 2174e6300d1a..a6604f7c707e 100644 --- a/crates/polars-plan/src/plans/optimizer/slice_pushdown_lp.rs +++ b/crates/polars-plan/src/plans/optimizer/slice_pushdown_lp.rs @@ -347,7 +347,7 @@ impl SlicePushDown { (Sort {input, by_column, slice, sort_options}, Some(state)) => { // The slice argument on Sort should be inserted by slice pushdown, // so it shouldn't exist yet (or be idempotently the same). - let new_slice = Some((state.offset, state.len as usize)); + let new_slice = Some((state.offset, state.len as usize, None)); assert!(slice.is_none() || slice == new_slice); // first restart optimization in inputs and get the updated LP diff --git a/crates/polars-stream/src/physical_plan/fmt.rs b/crates/polars-stream/src/physical_plan/fmt.rs index 4a99c8f603fb..6e5e056bb608 100644 --- a/crates/polars-stream/src/physical_plan/fmt.rs +++ b/crates/polars-stream/src/physical_plan/fmt.rs @@ -368,6 +368,7 @@ fn visualize_plan_rec( by_column, reverse, nulls_last: _, + dyn_pred: _, } => { let name = if reverse.iter().all(|r| *r) { "bottom-k" diff --git a/crates/polars-stream/src/physical_plan/lower_expr.rs b/crates/polars-stream/src/physical_plan/lower_expr.rs index d44f201305c5..8306841d15c5 100644 --- a/crates/polars-stream/src/physical_plan/lower_expr.rs +++ b/crates/polars-stream/src/physical_plan/lower_expr.rs @@ -1846,6 +1846,7 @@ fn lower_exprs_with_ctx( nulls_last: vec![true; by_column.len()], reverse, by_column, + dyn_pred: None, }; let output_schema = ctx.phys_sm[data_stream.node].output_schema.clone(); let node_key = ctx.phys_sm.insert(PhysNode::new(output_schema, kind)); diff --git a/crates/polars-stream/src/physical_plan/lower_ir.rs b/crates/polars-stream/src/physical_plan/lower_ir.rs index 382b3818009c..2470870cda7d 100644 --- a/crates/polars-stream/src/physical_plan/lower_ir.rs +++ b/crates/polars-stream/src/physical_plan/lower_ir.rs @@ -439,14 +439,14 @@ pub fn lower_ir( slice, sort_options, } => { - let slice = *slice; + let slice = slice.clone(); let mut by_column = by_column.clone(); let mut sort_options = sort_options.clone(); let phys_input = lower_ir!(*input)?; // See if we can insert a top k. let mut limit = u64::MAX; - if let Some((0, l)) = slice { + if let Some((0, l, _)) = slice { limit = limit.min(l as u64); } #[allow(clippy::unnecessary_cast)] @@ -503,6 +503,7 @@ pub fn lower_ir( by_column: trans_by_column, reverse: sort_options.descending.iter().map(|x| !x).collect(), nulls_last: sort_options.nulls_last.clone(), + dyn_pred: slice.as_ref().and_then(|t| t.2.clone()), }, })); } @@ -512,7 +513,7 @@ pub fn lower_ir( kind: PhysNodeKind::Sort { input: stream, by_column, - slice, + slice: slice.as_ref().map(|t| (t.0, t.1)), sort_options, }, })); diff --git a/crates/polars-stream/src/physical_plan/mod.rs b/crates/polars-stream/src/physical_plan/mod.rs index e77d26a7148c..fd05e5d2ce46 100644 --- a/crates/polars-stream/src/physical_plan/mod.rs +++ b/crates/polars-stream/src/physical_plan/mod.rs @@ -15,7 +15,7 @@ use polars_plan::dsl::{ }; use polars_plan::plans::expr_ir::ExprIR; use polars_plan::plans::hive::HivePartitionsDf; -use polars_plan::plans::{AExpr, DataFrameUdf, IR}; +use polars_plan::plans::{AExpr, DataFrameUdf, DynamicPred, IR}; mod fmt; mod io; @@ -239,6 +239,7 @@ pub enum PhysNodeKind { by_column: Vec, reverse: Vec, nulls_last: Vec, + dyn_pred: Option, }, Repeat { diff --git a/crates/polars-stream/src/physical_plan/to_graph.rs b/crates/polars-stream/src/physical_plan/to_graph.rs index 1ace1716d377..d86027dcf300 100644 --- a/crates/polars-stream/src/physical_plan/to_graph.rs +++ b/crates/polars-stream/src/physical_plan/to_graph.rs @@ -552,7 +552,7 @@ fn to_graph_rec<'a>( let sort_node = lp_arena.add(IR::Sort { input: df_node, by_column: by_column.clone(), - slice: *slice, + slice: slice.map(|t| (t.0, t.1, None)), sort_options: sort_options.clone(), }); let executor = Mutex::new(create_physical_plan( @@ -582,6 +582,7 @@ fn to_graph_rec<'a>( by_column, reverse, nulls_last, + dyn_pred: _, } => { let input_key = to_graph_rec(input.node, ctx)?; let k_key = to_graph_rec(k.node, ctx)?; From 9c827b41d75734134222d87c1ed1ee2cf5586c08 Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Sun, 8 Feb 2026 14:41:55 +0100 Subject: [PATCH 06/15] more dispatch --- crates/polars-plan/src/plans/optimizer/mod.rs | 22 +++++++++---------- .../optimizer/predicate_pushdown/dynamic.rs | 9 ++++---- .../plans/optimizer/predicate_pushdown/mod.rs | 10 ++++++--- .../src/lazyframe/visitor/expr_nodes.rs | 3 +++ .../src/lazyframe/visitor/nodes.rs | 2 +- crates/polars-stream/src/nodes/top_k.rs | 4 ++++ .../src/physical_plan/to_graph.rs | 3 ++- 7 files changed, 32 insertions(+), 21 deletions(-) diff --git a/crates/polars-plan/src/plans/optimizer/mod.rs b/crates/polars-plan/src/plans/optimizer/mod.rs index 21a0e52e272d..ffb2f6b90040 100644 --- a/crates/polars-plan/src/plans/optimizer/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/mod.rs @@ -196,17 +196,7 @@ pub fn optimize( true }; - if run_pushdowns { - run_projection_predicate_pushdown( - root, - ir_arena, - expr_arena, - pushdown_maintain_errors, - &opt_flags, - )?; - } - - if opt_flags.slice_pushdown() { + if dbg!(opt_flags.slice_pushdown()) { let mut slice_pushdown_opt = SlicePushDown::new( // We don't maintain errors on slice as the behavior is much more predictable that way. // @@ -223,6 +213,16 @@ pub fn optimize( rules.push(Box::new(slice_pushdown_opt)); } + if dbg!(run_pushdowns) { + run_projection_predicate_pushdown( + root, + ir_arena, + expr_arena, + pushdown_maintain_errors, + &opt_flags, + )?; + } + if opt_flags.fast_projection() { rules.push(Box::new(SimpleProjectionAndCollapse::new( opt_flags.eager(), diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index c24137a7a77b..a1409531de6e 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -69,11 +69,9 @@ impl DynamicPred { } pub fn evaluate(&self, columns: &[Column]) -> PolarsResult { - polars_ensure!(columns.len() > 1, ComputeError: "expected at least 1 argument"); - let h = columns[0].len(); - if self.inner.is_set.load(std::sync::atomic::Ordering::Relaxed) { + let out = if self.inner.is_set.load(std::sync::atomic::Ordering::Relaxed) { let guard = self.inner.pred.read().unwrap(); let dyn_func = guard.as_ref().unwrap(); dyn_func(columns) @@ -82,9 +80,10 @@ impl DynamicPred { Ok(Column::Scalar(ScalarColumn::new( columns[0].name().clone(), s, - h, + 1, ))) - } + }; + dbg!(out) } } diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs index f46eb7f55610..fd03fb6aa03a 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs @@ -578,16 +578,20 @@ impl PredicatePushDown { Sort { input, by_column, - slice, + mut slice, sort_options, } => { + dbg!(&slice); if let Some((offset, len, None)) = slice && by_column.len() == 1 { let n = by_column[0].node(); if let AExpr::Column(_) = expr_arena.get(n) { - let (node, pred) = new_dynamic_pred(n, expr_arena); - // todo + let (dyn_pred_node, pred) = new_dynamic_pred(n, expr_arena); + slice = Some((offset, len, Some(pred))); + + let predicate = ExprIR::from_node(dyn_pred_node, expr_arena); + insert_predicate_dedup(&mut acc_predicates, &predicate, expr_arena); } } diff --git a/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs b/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs index 3f6c5b83b08a..b1180d3cf88b 100644 --- a/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs +++ b/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs @@ -1429,6 +1429,9 @@ pub(crate) fn into_py(py: Python<'_>, expr: &AExpr) -> PyResult> { IRFunctionExpr::RowDecode(..) => { return Err(PyNotImplementedError::new_err("row_decode")); }, + IRFunctionExpr::DynamicExpr { .. } => { + return Err(PyNotImplementedError::new_err("dyn_expr")); + }, }?, options: py.None(), } diff --git a/crates/polars-python/src/lazyframe/visitor/nodes.rs b/crates/polars-python/src/lazyframe/visitor/nodes.rs index 332365d08ff8..3f98593da1fb 100644 --- a/crates/polars-python/src/lazyframe/visitor/nodes.rs +++ b/crates/polars-python/src/lazyframe/visitor/nodes.rs @@ -492,7 +492,7 @@ pub(crate) fn into_py(py: Python<'_>, plan: &IR) -> PyResult> { sort_options.nulls_last.clone(), sort_options.descending.clone(), ), - slice: *slice, + slice: slice.as_ref().map(|t| (t.0, t.1)), } .into_py_any(py), IR::Cache { input, id } => Cache { diff --git a/crates/polars-stream/src/nodes/top_k.rs b/crates/polars-stream/src/nodes/top_k.rs index 7f149f0440dd..61e6cbc17705 100644 --- a/crates/polars-stream/src/nodes/top_k.rs +++ b/crates/polars-stream/src/nodes/top_k.rs @@ -7,6 +7,7 @@ use polars_core::prelude::*; use polars_core::schema::Schema; use polars_core::utils::accumulate_dataframes_vertical; use polars_core::with_match_physical_numeric_polars_type; +use polars_plan::plans::DynamicPred; use polars_utils::IdxSize; use polars_utils::priority::Priority; use polars_utils::sort::ReorderWithNulls; @@ -399,6 +400,7 @@ pub struct TopKNode { key_schema: Arc, key_selectors: Vec, state: TopKState, + dyn_pred: Option, } impl TopKNode { @@ -408,6 +410,7 @@ impl TopKNode { nulls_last: Vec, key_schema: Arc, key_selectors: Vec, + dyn_pred: Option, ) -> Self { Self { reverse, @@ -415,6 +418,7 @@ impl TopKNode { key_schema, key_selectors, state: TopKState::WaitingForK(InMemorySinkNode::new(k_schema)), + dyn_pred, } } } diff --git a/crates/polars-stream/src/physical_plan/to_graph.rs b/crates/polars-stream/src/physical_plan/to_graph.rs index d86027dcf300..3e97c4269389 100644 --- a/crates/polars-stream/src/physical_plan/to_graph.rs +++ b/crates/polars-stream/src/physical_plan/to_graph.rs @@ -582,7 +582,7 @@ fn to_graph_rec<'a>( by_column, reverse, nulls_last, - dyn_pred: _, + dyn_pred, } => { let input_key = to_graph_rec(input.node, ctx)?; let k_key = to_graph_rec(k.node, ctx)?; @@ -603,6 +603,7 @@ fn to_graph_rec<'a>( nulls_last.clone(), key_schema, key_selectors, + dyn_pred, ), [(input_key, input.port), (k_key, k.port)], ) From 6a51b424fa89a7c4654d39f899e20bb13dc56129 Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Mon, 9 Feb 2026 13:59:25 +0100 Subject: [PATCH 07/15] remove dbg! --- crates/polars-plan/src/plans/optimizer/mod.rs | 4 ++-- .../src/plans/optimizer/predicate_pushdown/dynamic.rs | 7 ++++--- .../src/plans/optimizer/predicate_pushdown/mod.rs | 1 - 3 files changed, 6 insertions(+), 6 deletions(-) diff --git a/crates/polars-plan/src/plans/optimizer/mod.rs b/crates/polars-plan/src/plans/optimizer/mod.rs index ffb2f6b90040..cd1c623a49bd 100644 --- a/crates/polars-plan/src/plans/optimizer/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/mod.rs @@ -196,7 +196,7 @@ pub fn optimize( true }; - if dbg!(opt_flags.slice_pushdown()) { + if opt_flags.slice_pushdown() { let mut slice_pushdown_opt = SlicePushDown::new( // We don't maintain errors on slice as the behavior is much more predictable that way. // @@ -213,7 +213,7 @@ pub fn optimize( rules.push(Box::new(slice_pushdown_opt)); } - if dbg!(run_pushdowns) { + if run_pushdowns { run_projection_predicate_pushdown( root, ir_arena, diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index a1409531de6e..9d7259eb1a1b 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -71,7 +71,9 @@ impl DynamicPred { pub fn evaluate(&self, columns: &[Column]) -> PolarsResult { let h = columns[0].len(); - let out = if self.inner.is_set.load(std::sync::atomic::Ordering::Relaxed) { + // Can be relaxed, worst thing that can happen is that we read + // more data than strictly needed. + if self.inner.is_set.load(std::sync::atomic::Ordering::Relaxed) { let guard = self.inner.pred.read().unwrap(); let dyn_func = guard.as_ref().unwrap(); dyn_func(columns) @@ -82,8 +84,7 @@ impl DynamicPred { s, 1, ))) - }; - dbg!(out) + } } } diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs index fd03fb6aa03a..c5e76156b95f 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs @@ -581,7 +581,6 @@ impl PredicatePushDown { mut slice, sort_options, } => { - dbg!(&slice); if let Some((offset, len, None)) = slice && by_column.len() == 1 { From 55475ab2bb9d48c9e44f8dcc7545f518ac6c78c1 Mon Sep 17 00:00:00 2001 From: ritchie46 Date: Tue, 10 Feb 2026 11:32:42 +0100 Subject: [PATCH 08/15] use trait --- crates/polars-lazy/src/tests/optimization_checks.rs | 2 +- .../src/plans/optimizer/predicate_pushdown/dynamic.rs | 11 ++++++++--- crates/polars-stream/src/physical_plan/to_graph.rs | 2 +- 3 files changed, 10 insertions(+), 5 deletions(-) diff --git a/crates/polars-lazy/src/tests/optimization_checks.rs b/crates/polars-lazy/src/tests/optimization_checks.rs index f8ca17273274..211a05746b73 100644 --- a/crates/polars-lazy/src/tests/optimization_checks.rs +++ b/crates/polars-lazy/src/tests/optimization_checks.rs @@ -231,7 +231,7 @@ pub fn test_slice_pushdown_sort() -> PolarsResult<()> { assert!(lp_arena.iter(lp).all(|(_, lp)| { use IR::*; match lp { - Sort { slice, .. } => *slice == Some((1, 3)), + Sort { slice, .. } => *slice == Some((1, 3, None)), Slice { .. } => false, _ => true, } diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 9d7259eb1a1b..9d329726c64e 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -10,10 +10,15 @@ use serde::{Deserialize, Serialize}; use super::*; +pub trait DynamicExpr: Send + Sync { + // Invariant: Output Column must be of type `Boolean`. + fn evaluate(&self, columns: &[Column]) -> PolarsResult; +} + #[cfg_attr(feature = "ir_serde", derive(Serialize, Deserialize))] struct Inner { #[cfg_attr(feature = "ir_serde", serde(skip))] - pred: RwLock PolarsResult + Send + Sync>>>, + pred: RwLock>>, #[cfg_attr(feature = "ir_serde", serde(skip))] is_set: AtomicBool, id: UniqueId, @@ -58,7 +63,7 @@ impl DynamicPred { &self.inner.id } - pub fn set(&self, pred: Box PolarsResult + Send + Sync>) { + pub fn set(&self, pred: Box) { { let mut guard = self.inner.pred.write().unwrap(); *guard = Some(pred); @@ -76,7 +81,7 @@ impl DynamicPred { if self.inner.is_set.load(std::sync::atomic::Ordering::Relaxed) { let guard = self.inner.pred.read().unwrap(); let dyn_func = guard.as_ref().unwrap(); - dyn_func(columns) + dyn_func.evaluate(columns) } else { let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(true)); Ok(Column::Scalar(ScalarColumn::new( diff --git a/crates/polars-stream/src/physical_plan/to_graph.rs b/crates/polars-stream/src/physical_plan/to_graph.rs index 3e97c4269389..711fe7229646 100644 --- a/crates/polars-stream/src/physical_plan/to_graph.rs +++ b/crates/polars-stream/src/physical_plan/to_graph.rs @@ -603,7 +603,7 @@ fn to_graph_rec<'a>( nulls_last.clone(), key_schema, key_selectors, - dyn_pred, + dyn_pred.clone(), ), [(input_key, input.port), (k_key, k.port)], ) From edf3b0fbd957c56b1344dcc003a14c190e314f2f Mon Sep 17 00:00:00 2001 From: Orson Peters Date: Thu, 12 Feb 2026 10:59:40 +0100 Subject: [PATCH 09/15] Rename --- crates/polars-expr/src/dispatch/misc.rs | 2 +- crates/polars-expr/src/dispatch/mod.rs | 4 ++-- .../polars-plan/src/plans/aexpr/function_expr/mod.rs | 8 ++++---- .../src/plans/aexpr/function_expr/schema.rs | 2 +- crates/polars-plan/src/plans/conversion/ir_to_dsl.rs | 2 +- crates/polars-plan/src/plans/optimizer/mod.rs | 2 +- .../plans/optimizer/predicate_pushdown/dynamic.rs | 12 +++++++----- .../src/plans/optimizer/predicate_pushdown/mod.rs | 2 +- .../src/lazyframe/visitor/expr_nodes.rs | 4 ++-- crates/polars-stream/src/nodes/top_k.rs | 12 +++++++++++- 10 files changed, 31 insertions(+), 19 deletions(-) diff --git a/crates/polars-expr/src/dispatch/misc.rs b/crates/polars-expr/src/dispatch/misc.rs index 03527ea36455..bd5446d73db5 100644 --- a/crates/polars-expr/src/dispatch/misc.rs +++ b/crates/polars-expr/src/dispatch/misc.rs @@ -1055,6 +1055,6 @@ pub fn repeat(args: &[Column]) -> PolarsResult { Ok(c.new_from_index(0, n)) } -pub fn dynamic_expr(columns: &[Column], pred: &DynamicPred) -> PolarsResult { +pub fn dynamic_pred(columns: &[Column], pred: &DynamicPred) -> PolarsResult { pred.evaluate(columns) } diff --git a/crates/polars-expr/src/dispatch/mod.rs b/crates/polars-expr/src/dispatch/mod.rs index 90a4dc68f84d..3f68ad8ed476 100644 --- a/crates/polars-expr/src/dispatch/mod.rs +++ b/crates/polars-expr/src/dispatch/mod.rs @@ -527,8 +527,8 @@ pub fn function_expr_to_udf(func: IRFunctionExpr) -> SpecialEq { map_as_slice!(misc::row_decode, fs.clone(), variants.clone()) }, - F::DynamicExpr { pred } => { - map_as_slice!(misc::dynamic_expr, &pred) + F::DynamicPred { pred } => { + map_as_slice!(misc::dynamic_pred, &pred) }, } } diff --git a/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs b/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs index bbab33ed4ea2..e6fff3bfd06a 100644 --- a/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs +++ b/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs @@ -388,7 +388,7 @@ pub enum IRFunctionExpr { RowEncode(Vec, RowEncodingVariant), #[cfg(feature = "dtype-struct")] RowDecode(Vec, RowEncodingVariant), - DynamicExpr { + DynamicPred { pred: DynamicPred, }, } @@ -694,7 +694,7 @@ impl Hash for IRFunctionExpr { fs.hash(state); variants.hash(state); }, - DynamicExpr { pred } => { + DynamicPred { pred } => { pred.id().hash(state); }, } @@ -911,7 +911,7 @@ impl Display for IRFunctionExpr { RowEncode(..) => "row_encode", #[cfg(feature = "dtype-struct")] RowDecode(..) => "row_decode", - DynamicExpr { pred } => "dynamic_predicate", + DynamicPred { pred } => "dynamic_predicate", }; write!(f, "{s}") } @@ -1241,7 +1241,7 @@ impl IRFunctionExpr { F::RowEncode(..) => FunctionOptions::elementwise(), #[cfg(feature = "dtype-struct")] F::RowDecode(..) => FunctionOptions::elementwise(), - F::DynamicExpr { .. } => FunctionOptions::elementwise(), + F::DynamicPred { .. } => FunctionOptions::elementwise(), } } } diff --git a/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs b/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs index 49a91adfa986..d233b5d0f562 100644 --- a/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs +++ b/crates/polars-plan/src/plans/aexpr/function_expr/schema.rs @@ -451,7 +451,7 @@ impl IRFunctionExpr { }), #[cfg(feature = "dtype-struct")] RowDecode(fields, _) => mapper.with_dtype(DataType::Struct(fields.to_vec())), - DynamicExpr { .. } => mapper.with_dtype(DataType::Boolean), + DynamicPred { .. } => mapper.with_dtype(DataType::Boolean), } } diff --git a/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs b/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs index 9121db6ce863..f3cbdfccf763 100644 --- a/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs +++ b/crates/polars-plan/src/plans/conversion/ir_to_dsl.rs @@ -1177,7 +1177,7 @@ pub fn ir_function_to_dsl(input: Vec, function: IRFunctionExpr) -> Expr { fs.into_iter().map(|f| (f.name, f.dtype.into())).collect(), v, ), - IF::DynamicExpr { pred } => { + IF::DynamicPred { pred } => { return Expr::Display { inputs: input, fmt_str: Box::new(format_pl_smallstr!("{pred:?}")), diff --git a/crates/polars-plan/src/plans/optimizer/mod.rs b/crates/polars-plan/src/plans/optimizer/mod.rs index cd1c623a49bd..089603dfa2b4 100644 --- a/crates/polars-plan/src/plans/optimizer/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/mod.rs @@ -34,7 +34,7 @@ pub use cse::NaiveExprMerger; use delay_rechunk::DelayRechunk; pub use expand_datasets::ExpandedDataset; use polars_core::config::verbose; -pub use predicate_pushdown::{DynamicPred, PredicatePushDown}; +pub use predicate_pushdown::{DynamicPred, PredicateExpr, PredicatePushDown}; pub use projection_pushdown::ProjectionPushDown; pub use simplify_expr::{SimplifyBooleanRule, SimplifyExprRule}; use slice_pushdown_lp::SlicePushDown; diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 9d329726c64e..94c4e0fe84fe 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -1,3 +1,4 @@ +use std::any::Any; use std::fmt::{Debug, Formatter}; use std::hash::Hash; use std::sync::RwLock; @@ -10,15 +11,16 @@ use serde::{Deserialize, Serialize}; use super::*; -pub trait DynamicExpr: Send + Sync { - // Invariant: Output Column must be of type `Boolean`. +pub trait PredicateExpr: Send + Sync + Any { + // Invariant: output column must be of type `Boolean`. If true a value is + // included, if false it is filtered out. fn evaluate(&self, columns: &[Column]) -> PolarsResult; } #[cfg_attr(feature = "ir_serde", derive(Serialize, Deserialize))] struct Inner { #[cfg_attr(feature = "ir_serde", serde(skip))] - pred: RwLock>>, + pred: RwLock>>, #[cfg_attr(feature = "ir_serde", serde(skip))] is_set: AtomicBool, id: UniqueId, @@ -63,7 +65,7 @@ impl DynamicPred { &self.inner.id } - pub fn set(&self, pred: Box) { + pub fn set(&self, pred: Arc) { { let mut guard = self.inner.pred.write().unwrap(); *guard = Some(pred); @@ -95,7 +97,7 @@ impl DynamicPred { pub fn new_dynamic_pred(node: Node, arena: &mut Arena) -> (Node, DynamicPred) { let pred = DynamicPred::new(); - let function = IRFunctionExpr::DynamicExpr { pred: pred.clone() }; + let function = IRFunctionExpr::DynamicPred { pred: pred.clone() }; let options = function.function_options(); let aexpr = AExpr::Function { input: vec![ExprIR::from_node(node, arena)], diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs index c5e76156b95f..b31aa679d964 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs @@ -4,7 +4,7 @@ mod join; mod keys; mod utils; -pub use dynamic::DynamicPred; +pub use dynamic::{DynamicPred, PredicateExpr}; use polars_core::datatypes::PlHashMap; use polars_core::prelude::*; use polars_utils::idx_vec::UnitVec; diff --git a/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs b/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs index b1180d3cf88b..dc83f8e1a670 100644 --- a/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs +++ b/crates/polars-python/src/lazyframe/visitor/expr_nodes.rs @@ -1429,8 +1429,8 @@ pub(crate) fn into_py(py: Python<'_>, expr: &AExpr) -> PyResult> { IRFunctionExpr::RowDecode(..) => { return Err(PyNotImplementedError::new_err("row_decode")); }, - IRFunctionExpr::DynamicExpr { .. } => { - return Err(PyNotImplementedError::new_err("dyn_expr")); + IRFunctionExpr::DynamicPred { .. } => { + return Err(PyNotImplementedError::new_err("dynamic_pred")); }, }?, options: py.None(), diff --git a/crates/polars-stream/src/nodes/top_k.rs b/crates/polars-stream/src/nodes/top_k.rs index 61e6cbc17705..81213e16934b 100644 --- a/crates/polars-stream/src/nodes/top_k.rs +++ b/crates/polars-stream/src/nodes/top_k.rs @@ -7,7 +7,7 @@ use polars_core::prelude::*; use polars_core::schema::Schema; use polars_core::utils::accumulate_dataframes_vertical; use polars_core::with_match_physical_numeric_polars_type; -use polars_plan::plans::DynamicPred; +use polars_plan::plans::{DynamicPred, PredicateExpr}; use polars_utils::IdxSize; use polars_utils::priority::Priority; use polars_utils::sort::ReorderWithNulls; @@ -412,6 +412,16 @@ impl TopKNode { key_selectors: Vec, dyn_pred: Option, ) -> Self { + if let Some(p) = &dyn_pred { + struct Test; + impl PredicateExpr for Test { + fn evaluate(&self, columns: &[Column]) -> PolarsResult { + unimplemented!() + } + } + + p.set(Arc::new(Test)); + } Self { reverse, nulls_last, From 842f8eddcb0d7a79378272edfd78d4025f2771d8 Mon Sep 17 00:00:00 2001 From: Orson Peters Date: Thu, 12 Feb 2026 16:08:31 +0100 Subject: [PATCH 10/15] Implement in streaming engine --- crates/polars-plan/src/plans/optimizer/mod.rs | 2 +- .../optimizer/predicate_pushdown/dynamic.rs | 40 +++-- .../plans/optimizer/predicate_pushdown/mod.rs | 2 +- crates/polars-stream/src/nodes/top_k.rs | 160 +++++++++++++++--- crates/polars-utils/src/sort.rs | 10 ++ 5 files changed, 171 insertions(+), 43 deletions(-) diff --git a/crates/polars-plan/src/plans/optimizer/mod.rs b/crates/polars-plan/src/plans/optimizer/mod.rs index 089603dfa2b4..903a43127cd8 100644 --- a/crates/polars-plan/src/plans/optimizer/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/mod.rs @@ -34,7 +34,7 @@ pub use cse::NaiveExprMerger; use delay_rechunk::DelayRechunk; pub use expand_datasets::ExpandedDataset; use polars_core::config::verbose; -pub use predicate_pushdown::{DynamicPred, PredicateExpr, PredicatePushDown}; +pub use predicate_pushdown::{DynamicPred, PredicateExpr, PredicatePushDown, TrivialPredicateExpr}; pub use projection_pushdown::ProjectionPushDown; pub use simplify_expr::{SimplifyBooleanRule, SimplifyExprRule}; use slice_pushdown_lp::SlicePushDown; diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 94c4e0fe84fe..687eb77b929b 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -2,7 +2,7 @@ use std::any::Any; use std::fmt::{Debug, Formatter}; use std::hash::Hash; use std::sync::RwLock; -use std::sync::atomic::AtomicBool; +use std::sync::atomic::{AtomicBool, Ordering}; use polars_core::frame::column::ScalarColumn; use polars_utils::unique_id::UniqueId; @@ -13,8 +13,17 @@ use super::*; pub trait PredicateExpr: Send + Sync + Any { // Invariant: output column must be of type `Boolean`. If true a value is - // included, if false it is filtered out. - fn evaluate(&self, columns: &[Column]) -> PolarsResult; + // included, if false it is filtered out. If None is returned it is assumed + // all values are needed. + fn evaluate(&self, columns: &[Column]) -> PolarsResult>; +} + +pub struct TrivialPredicateExpr; + +impl PredicateExpr for TrivialPredicateExpr { + fn evaluate(&self, columns: &[Column]) -> PolarsResult> { + Ok(None) + } } #[cfg_attr(feature = "ir_serde", derive(Serialize, Deserialize))] @@ -72,26 +81,25 @@ impl DynamicPred { } self.inner .is_set - .store(true, std::sync::atomic::Ordering::Release); + .store(true, Ordering::Release); } pub fn evaluate(&self, columns: &[Column]) -> PolarsResult { let h = columns[0].len(); - - // Can be relaxed, worst thing that can happen is that we read - // more data than strictly needed. - if self.inner.is_set.load(std::sync::atomic::Ordering::Relaxed) { + if self.inner.is_set.load(Ordering::Acquire) { let guard = self.inner.pred.read().unwrap(); let dyn_func = guard.as_ref().unwrap(); - dyn_func.evaluate(columns) - } else { - let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(true)); - Ok(Column::Scalar(ScalarColumn::new( - columns[0].name().clone(), - s, - 1, - ))) + if let Some(pred) = dyn_func.evaluate(columns)? { + return Ok(pred); + } } + + let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(true)); + Ok(Column::Scalar(ScalarColumn::new( + columns[0].name().clone(), + s, + 1, + ))) } } diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs index b31aa679d964..dc0a9521d188 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/mod.rs @@ -4,7 +4,7 @@ mod join; mod keys; mod utils; -pub use dynamic::{DynamicPred, PredicateExpr}; +pub use dynamic::{DynamicPred, PredicateExpr, TrivialPredicateExpr}; use polars_core::datatypes::PlHashMap; use polars_core::prelude::*; use polars_utils::idx_vec::UnitVec; diff --git a/crates/polars-stream/src/nodes/top_k.rs b/crates/polars-stream/src/nodes/top_k.rs index 81213e16934b..6c7d19f727ca 100644 --- a/crates/polars-stream/src/nodes/top_k.rs +++ b/crates/polars-stream/src/nodes/top_k.rs @@ -2,12 +2,13 @@ use std::any::Any; use std::collections::BinaryHeap; use std::sync::Arc; +use parking_lot::RwLock; use polars_core::prelude::row_encode::_get_rows_encoded; use polars_core::prelude::*; use polars_core::schema::Schema; use polars_core::utils::accumulate_dataframes_vertical; use polars_core::with_match_physical_numeric_polars_type; -use polars_plan::plans::{DynamicPred, PredicateExpr}; +use polars_plan::plans::{DynamicPred, PredicateExpr, TrivialPredicateExpr}; use polars_utils::IdxSize; use polars_utils::priority::Priority; use polars_utils::sort::ReorderWithNulls; @@ -61,17 +62,18 @@ impl DfSubset { } } -pub struct BottomKWithPayload

{ +struct BottomKWithPayload

{ k: usize, heap: BinaryHeap>, df_subsets: SlotMap, row_idxs: SlotMap, to_prune: SecondaryMap, gather_idxs: Vec, + shared_optimum: Arc>>, } impl BottomKWithPayload

{ - pub fn new(k: usize) -> Self { + pub fn new(k: usize, shared_optimum: Arc>>) -> Self { Self { k, heap: BinaryHeap::with_capacity(k + 1), @@ -79,6 +81,7 @@ impl BottomKWithPayload

{ row_idxs: SlotMap::with_key(), to_prune: SecondaryMap::new(), gather_idxs: Vec::new(), + shared_optimum, } } @@ -87,6 +90,7 @@ impl BottomKWithPayload

{ df: DataFrame, keys: impl IntoIterator, is_less: impl Fn(&Q, &P) -> bool, + is_less_owned: impl Fn(&P, &P) -> bool, to_owned: impl Fn(Q) -> P, ) { let dfs_key = self.df_subsets.insert(DfSubset { @@ -95,16 +99,29 @@ impl BottomKWithPayload

{ subset_len: 0, }); + let mut new_optimum = false; for (row_idx, key) in keys.into_iter().enumerate() { - self.add_one( + new_optimum |= self.add_one( dfs_key, row_idx.try_into().unwrap(), key, &is_less, &to_owned, - ) + ); } self.prune(); + + if new_optimum { + let new_shared_opt = if let Some(v) = self.shared_optimum.read().clone() { + is_less_owned(self.peek_optimum().unwrap(), &v) + } else { + true + }; + + if new_shared_opt { + *self.shared_optimum.write() = self.peek_optimum().cloned(); + } + } } fn add_one( @@ -114,16 +131,19 @@ impl BottomKWithPayload

{ key: Q, is_less: impl Fn(&Q, &P) -> bool, to_owned: impl Fn(Q) -> P, - ) { + ) -> bool { // We use a max-heap for our bottom k. This means the top element in our heap (peek()) // is the first to be replaced. - if self.heap.len() < self.k || is_less(&key, &self.heap.peek().unwrap().0) { + let mut new_optimum = false; + if self.heap.len() < self.k || is_less(&key, self.peek_optimum().unwrap()) { let row_idx_key = self.row_idxs.insert(row_idx); let df_subset = &mut self.df_subsets[dfs_key]; df_subset.subset_len += 1; df_subset.rows.push(row_idx_key); + let opt = self.heap.peek().map(|p| p.1); self.heap .push(Priority(to_owned(key), (dfs_key, row_idx_key))); + new_optimum = opt != self.heap.peek().map(|p| p.1); } if self.heap.len() > self.k { @@ -135,6 +155,8 @@ impl BottomKWithPayload

{ self.to_prune.insert(dfs_key, ()); } } + + new_optimum } pub fn prune(&mut self) { @@ -190,10 +212,15 @@ impl BottomKWithPayload

{ self.to_prune.clear(); Some(ret.unwrap()) } + + fn peek_optimum(&self) -> Option<&P> { + self.heap.peek().map(|x| &x.0) + } } trait DfByKeyReducer: Any + Send + 'static { fn new_empty(&self) -> Box; + fn new_pred(&self) -> Arc; fn add(&mut self, df: DataFrame, keys: DataFrame); fn combine(&mut self, other: &dyn DfByKeyReducer); fn finalize(self: Box) -> Option; @@ -210,7 +237,7 @@ impl { fn new(k: usize) -> Self { Self { - inner: BottomKWithPayload::new(k), + inner: BottomKWithPayload::new(k, Arc::default()), } } } @@ -220,7 +247,13 @@ impl DfByKeyR { fn new_empty(&self) -> Box { Box::new(Self { - inner: BottomKWithPayload::new(self.inner.k), + inner: BottomKWithPayload::new(self.inner.k, self.inner.shared_optimum.clone()), + }) + } + + fn new_pred(&self) -> Arc { + Arc::new(PrimitiveBottomKPredicate:: { + shared_optimum: self.inner.shared_optimum.clone() }) } @@ -234,6 +267,7 @@ impl DfByKeyR .iter() .map(|opt_x| ReorderWithNulls(opt_x.map(TotalOrdWrap))), |l, r| l < r, + |l, r| l < r, |x| x, ); } @@ -248,6 +282,35 @@ impl DfByKeyR } } +struct PrimitiveBottomKPredicate { + shared_optimum: Arc>, REVERSE, NULLS_LAST>>>>, +} + +impl PredicateExpr + for PrimitiveBottomKPredicate +{ + fn evaluate(&self, columns: &[Column]) -> PolarsResult> { + let Some(v) = self.shared_optimum.read().clone() else { + return Ok(None); + }; + + if columns[0].dtype().is_null() || matches!(columns[0], Column::Scalar(_)) { + let cv = columns[0].get(0)?.null_to_none().map(|v| TotalOrdWrap(v.try_extract().unwrap())); + let keep = ReorderWithNulls(cv) < v; + let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(keep)); + Ok(Some(Column::new_scalar(PlSmallStr::EMPTY, s, columns[0].len()))) + } else { + let keys = columns[0].as_materialized_series(); + let key_ca: &ChunkedArray = keys.as_phys_any().downcast_ref().unwrap(); + let pred: BooleanChunked = key_ca + .iter() + .map(|opt_x| ReorderWithNulls(opt_x.map(TotalOrdWrap)) < v) + .collect_ca(PlSmallStr::EMPTY); + Ok(Some(Column::from(pred.into_series()))) + } + } +} + struct BinaryBottomK { inner: BottomKWithPayload, REVERSE, NULLS_LAST>>, } @@ -255,7 +318,7 @@ struct BinaryBottomK { impl BinaryBottomK { fn new(k: usize) -> Self { Self { - inner: BottomKWithPayload::new(k), + inner: BottomKWithPayload::new(k, Arc::default()), } } } @@ -265,7 +328,13 @@ impl DfByKeyReducer { fn new_empty(&self) -> Box { Box::new(Self { - inner: BottomKWithPayload::new(self.inner.k), + inner: BottomKWithPayload::new(self.inner.k, self.inner.shared_optimum.clone()), + }) + } + + fn new_pred(&self) -> Arc { + Arc::new(BinaryBottomKPredicate { + shared_optimum: self.inner.shared_optimum.clone() }) } @@ -274,14 +343,15 @@ impl DfByKeyReducer let key_ca = if let Ok(ca_str) = keys[0].str() { ca_str.as_binary() } else { - df[0].binary().unwrap().clone() + keys[0].binary().unwrap().clone() }; self.inner.add_df( df, key_ca .iter() .map(ReorderWithNulls::<_, REVERSE, NULLS_LAST>), - |l, r| l < &ReorderWithNulls(r.0.as_deref()), + |l, r| l < &r.as_deref(), + |l, r| l < r, |x| ReorderWithNulls(x.0.map(<[u8]>::to_vec)), ); } @@ -296,6 +366,47 @@ impl DfByKeyReducer } } +struct BinaryBottomKPredicate { + shared_optimum: Arc, REVERSE, NULLS_LAST>>>>, +} + +impl PredicateExpr + for BinaryBottomKPredicate +{ + fn evaluate(&self, columns: &[Column]) -> PolarsResult> { + let Some(v) = self.shared_optimum.read().clone() else { + return Ok(None); + }; + + if columns[0].dtype().is_null() || matches!(columns[0], Column::Scalar(_)) { + let scalar = columns[0].get(0)?; + let cv = match &scalar { + AnyValue::Null => None, + AnyValue::String(s) => Some(s.as_bytes()), + AnyValue::StringOwned(s) => Some(s.as_bytes()), + AnyValue::Binary(b) => Some(*b), + AnyValue::BinaryOwned(b) => Some(b.as_slice()), + _ => unreachable!(), + }; + let keep = ReorderWithNulls(cv) < v.as_deref(); + let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(keep)); + Ok(Some(Column::new_scalar(PlSmallStr::EMPTY, s, columns[0].len()))) + } else { + let keys = columns[0].as_materialized_series(); + let key_ca = if let Ok(ca_str) = keys.str() { + ca_str.as_binary() + } else { + keys.binary().unwrap().clone() + }; + let pred: BooleanChunked = key_ca + .iter() + .map(|opt_x| ReorderWithNulls(opt_x) < v.as_deref()) + .collect_ca(PlSmallStr::EMPTY); + Ok(Some(Column::from(pred.into_series()))) + } + } +} + struct RowEncodedBottomK { inner: BottomKWithPayload>, reverse: Vec, @@ -305,7 +416,7 @@ struct RowEncodedBottomK { impl RowEncodedBottomK { fn new(k: usize, reverse: Vec, nulls_last: Vec) -> Self { Self { - inner: BottomKWithPayload::new(k), + inner: BottomKWithPayload::new(k, Arc::default()), reverse, nulls_last, } @@ -315,12 +426,17 @@ impl RowEncodedBottomK { impl DfByKeyReducer for RowEncodedBottomK { fn new_empty(&self) -> Box { Box::new(Self { - inner: BottomKWithPayload::new(self.inner.k), + inner: BottomKWithPayload::new(self.inner.k, self.inner.shared_optimum.clone()), reverse: self.reverse.clone(), nulls_last: self.nulls_last.clone(), }) } + fn new_pred(&self) -> Arc { + // Not implemented for row-encoded keys. + Arc::new(TrivialPredicateExpr) + } + fn add(&mut self, df: DataFrame, keys: DataFrame) { let keys_encoded = _get_rows_encoded(keys.columns(), &self.reverse, &self.nulls_last) .unwrap() @@ -329,6 +445,7 @@ impl DfByKeyReducer for RowEncodedBottomK { df, keys_encoded.values_iter(), |l, r| *l < r.as_slice(), + |l, r| l < r, |x| x.to_vec(), ); } @@ -412,16 +529,6 @@ impl TopKNode { key_selectors: Vec, dyn_pred: Option, ) -> Self { - if let Some(p) = &dyn_pred { - struct Test; - impl PredicateExpr for Test { - fn evaluate(&self, columns: &[Column]) -> PolarsResult { - unimplemented!() - } - } - - p.set(Arc::new(Test)); - } Self { reverse, nulls_last, @@ -468,6 +575,9 @@ impl ComputeNode for TopKNode { if k > 0 { let reducer = new_top_k_reducer(k, &self.reverse, &self.nulls_last, &self.key_schema); + if let Some(dyn_pred) = &self.dyn_pred { + dyn_pred.set(reducer.new_pred()); + } let reducers = (0..state.num_pipelines) .map(|_| reducer.new_empty()) .collect(); diff --git a/crates/polars-utils/src/sort.rs b/crates/polars-utils/src/sort.rs index 9878ab358015..4e739599021a 100644 --- a/crates/polars-utils/src/sort.rs +++ b/crates/polars-utils/src/sort.rs @@ -1,5 +1,6 @@ use std::cmp::Ordering; use std::mem::MaybeUninit; +use std::ops::Deref; use num_traits::FromPrimitive; @@ -57,6 +58,15 @@ where #[repr(transparent)] pub struct ReorderWithNulls(pub Option); +impl ReorderWithNulls { + pub fn as_deref(&self) -> ReorderWithNulls<&::Target, DESCENDING, NULLS_LAST> + where T: Deref + { + let x = self.0.as_deref(); + ReorderWithNulls(x) + } +} + impl PartialOrd for ReorderWithNulls { From 55dfb362fa7cf0228064261d2a3eb323d7ba2392 Mon Sep 17 00:00:00 2001 From: Orson Peters Date: Thu, 12 Feb 2026 16:08:42 +0100 Subject: [PATCH 11/15] Fmt --- .../optimizer/predicate_pushdown/dynamic.rs | 4 +- crates/polars-stream/src/nodes/top_k.rs | 40 +++++++++++++------ crates/polars-utils/src/sort.rs | 7 +++- 3 files changed, 33 insertions(+), 18 deletions(-) diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index 687eb77b929b..dcaee9fefa7f 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -79,9 +79,7 @@ impl DynamicPred { let mut guard = self.inner.pred.write().unwrap(); *guard = Some(pred); } - self.inner - .is_set - .store(true, Ordering::Release); + self.inner.is_set.store(true, Ordering::Release); } pub fn evaluate(&self, columns: &[Column]) -> PolarsResult { diff --git a/crates/polars-stream/src/nodes/top_k.rs b/crates/polars-stream/src/nodes/top_k.rs index 6c7d19f727ca..d519af276b21 100644 --- a/crates/polars-stream/src/nodes/top_k.rs +++ b/crates/polars-stream/src/nodes/top_k.rs @@ -117,7 +117,7 @@ impl BottomKWithPayload

{ } else { true }; - + if new_shared_opt { *self.shared_optimum.write() = self.peek_optimum().cloned(); } @@ -155,7 +155,7 @@ impl BottomKWithPayload

{ self.to_prune.insert(dfs_key, ()); } } - + new_optimum } @@ -212,7 +212,7 @@ impl BottomKWithPayload

{ self.to_prune.clear(); Some(ret.unwrap()) } - + fn peek_optimum(&self) -> Option<&P> { self.heap.peek().map(|x| &x.0) } @@ -250,10 +250,10 @@ impl DfByKeyR inner: BottomKWithPayload::new(self.inner.k, self.inner.shared_optimum.clone()), }) } - + fn new_pred(&self) -> Arc { Arc::new(PrimitiveBottomKPredicate:: { - shared_optimum: self.inner.shared_optimum.clone() + shared_optimum: self.inner.shared_optimum.clone(), }) } @@ -282,8 +282,11 @@ impl DfByKeyR } } -struct PrimitiveBottomKPredicate { - shared_optimum: Arc>, REVERSE, NULLS_LAST>>>>, +struct PrimitiveBottomKPredicate +{ + shared_optimum: Arc< + RwLock>, REVERSE, NULLS_LAST>>>, + >, } impl PredicateExpr @@ -293,12 +296,19 @@ impl Predicat let Some(v) = self.shared_optimum.read().clone() else { return Ok(None); }; - + if columns[0].dtype().is_null() || matches!(columns[0], Column::Scalar(_)) { - let cv = columns[0].get(0)?.null_to_none().map(|v| TotalOrdWrap(v.try_extract().unwrap())); + let cv = columns[0] + .get(0)? + .null_to_none() + .map(|v| TotalOrdWrap(v.try_extract().unwrap())); let keep = ReorderWithNulls(cv) < v; let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(keep)); - Ok(Some(Column::new_scalar(PlSmallStr::EMPTY, s, columns[0].len()))) + Ok(Some(Column::new_scalar( + PlSmallStr::EMPTY, + s, + columns[0].len(), + ))) } else { let keys = columns[0].as_materialized_series(); let key_ca: &ChunkedArray = keys.as_phys_any().downcast_ref().unwrap(); @@ -334,7 +344,7 @@ impl DfByKeyReducer fn new_pred(&self) -> Arc { Arc::new(BinaryBottomKPredicate { - shared_optimum: self.inner.shared_optimum.clone() + shared_optimum: self.inner.shared_optimum.clone(), }) } @@ -377,7 +387,7 @@ impl PredicateExpr let Some(v) = self.shared_optimum.read().clone() else { return Ok(None); }; - + if columns[0].dtype().is_null() || matches!(columns[0], Column::Scalar(_)) { let scalar = columns[0].get(0)?; let cv = match &scalar { @@ -390,7 +400,11 @@ impl PredicateExpr }; let keep = ReorderWithNulls(cv) < v.as_deref(); let s = Scalar::new(DataType::Boolean, AnyValue::Boolean(keep)); - Ok(Some(Column::new_scalar(PlSmallStr::EMPTY, s, columns[0].len()))) + Ok(Some(Column::new_scalar( + PlSmallStr::EMPTY, + s, + columns[0].len(), + ))) } else { let keys = columns[0].as_materialized_series(); let key_ca = if let Ok(ca_str) = keys.str() { diff --git a/crates/polars-utils/src/sort.rs b/crates/polars-utils/src/sort.rs index 4e739599021a..d7a96ac43e79 100644 --- a/crates/polars-utils/src/sort.rs +++ b/crates/polars-utils/src/sort.rs @@ -58,9 +58,12 @@ where #[repr(transparent)] pub struct ReorderWithNulls(pub Option); -impl ReorderWithNulls { +impl + ReorderWithNulls +{ pub fn as_deref(&self) -> ReorderWithNulls<&::Target, DESCENDING, NULLS_LAST> - where T: Deref + where + T: Deref, { let x = self.0.as_deref(); ReorderWithNulls(x) From 8d9022f30e8b1e5f66fd97ceb9a7940a1df0ed7b Mon Sep 17 00:00:00 2001 From: Orson Peters Date: Tue, 17 Feb 2026 13:27:22 +0100 Subject: [PATCH 12/15] Fix tests --- crates/polars-plan/src/plans/aexpr/function_expr/mod.rs | 2 +- crates/polars-plan/src/plans/ir/format.rs | 4 ++-- .../src/plans/optimizer/predicate_pushdown/dynamic.rs | 7 +++---- crates/polars-stream/src/nodes/top_k.rs | 3 ++- 4 files changed, 8 insertions(+), 8 deletions(-) diff --git a/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs b/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs index e6fff3bfd06a..a0f11b293921 100644 --- a/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs +++ b/crates/polars-plan/src/plans/aexpr/function_expr/mod.rs @@ -911,7 +911,7 @@ impl Display for IRFunctionExpr { RowEncode(..) => "row_encode", #[cfg(feature = "dtype-struct")] RowDecode(..) => "row_decode", - DynamicPred { pred } => "dynamic_predicate", + DynamicPred { .. } => "dynamic_predicate", }; write!(f, "{s}") } diff --git a/crates/polars-plan/src/plans/ir/format.rs b/crates/polars-plan/src/plans/ir/format.rs index 89adb3a7c6c7..ad8a206c382b 100644 --- a/crates/polars-plan/src/plans/ir/format.rs +++ b/crates/polars-plan/src/plans/ir/format.rs @@ -853,9 +853,9 @@ pub fn write_ir_non_recursive( let mut comma = false; if let Some((o, l, dyn_pred)) = slice { if let Some(dyn_pred) = &dyn_pred { - write!(f, "slice: ({o}, {l})")?; - } else { write!(f, "slice: ({o}, {l}, {dyn_pred:?})")?; + } else { + write!(f, "slice: ({o}, {l})")?; } comma = true; } diff --git a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs index dcaee9fefa7f..13c31cb791f4 100644 --- a/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs +++ b/crates/polars-plan/src/plans/optimizer/predicate_pushdown/dynamic.rs @@ -6,7 +6,7 @@ use std::sync::atomic::{AtomicBool, Ordering}; use polars_core::frame::column::ScalarColumn; use polars_utils::unique_id::UniqueId; -#[cfg(feature = "serde")] +#[cfg(feature = "ir_serde")] use serde::{Deserialize, Serialize}; use super::*; @@ -21,7 +21,7 @@ pub trait PredicateExpr: Send + Sync + Any { pub struct TrivialPredicateExpr; impl PredicateExpr for TrivialPredicateExpr { - fn evaluate(&self, columns: &[Column]) -> PolarsResult> { + fn evaluate(&self, _columns: &[Column]) -> PolarsResult> { Ok(None) } } @@ -83,7 +83,6 @@ impl DynamicPred { } pub fn evaluate(&self, columns: &[Column]) -> PolarsResult { - let h = columns[0].len(); if self.inner.is_set.load(Ordering::Acquire) { let guard = self.inner.pred.read().unwrap(); let dyn_func = guard.as_ref().unwrap(); @@ -96,7 +95,7 @@ impl DynamicPred { Ok(Column::Scalar(ScalarColumn::new( columns[0].name().clone(), s, - 1, + columns[0].len(), ))) } } diff --git a/crates/polars-stream/src/nodes/top_k.rs b/crates/polars-stream/src/nodes/top_k.rs index d519af276b21..d9e92cd73078 100644 --- a/crates/polars-stream/src/nodes/top_k.rs +++ b/crates/polars-stream/src/nodes/top_k.rs @@ -111,7 +111,7 @@ impl BottomKWithPayload

{ } self.prune(); - if new_optimum { + if new_optimum && self.heap.len() == self.k { let new_shared_opt = if let Some(v) = self.shared_optimum.read().clone() { is_less_owned(self.peek_optimum().unwrap(), &v) } else { @@ -284,6 +284,7 @@ impl DfByKeyR struct PrimitiveBottomKPredicate { + #[allow(clippy::type_complexity)] shared_optimum: Arc< RwLock>, REVERSE, NULLS_LAST>>>, >, From 2fd47e21011e0c0f930b005b1d6f06decccf5b6c Mon Sep 17 00:00:00 2001 From: Orson Peters Date: Tue, 17 Feb 2026 14:26:03 +0100 Subject: [PATCH 13/15] Fix rust test --- crates/polars-lazy/src/tests/optimization_checks.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/polars-lazy/src/tests/optimization_checks.rs b/crates/polars-lazy/src/tests/optimization_checks.rs index 211a05746b73..22221a9119d7 100644 --- a/crates/polars-lazy/src/tests/optimization_checks.rs +++ b/crates/polars-lazy/src/tests/optimization_checks.rs @@ -231,7 +231,7 @@ pub fn test_slice_pushdown_sort() -> PolarsResult<()> { assert!(lp_arena.iter(lp).all(|(_, lp)| { use IR::*; match lp { - Sort { slice, .. } => *slice == Some((1, 3, None)), + Sort { slice, .. } => matches!(slice, Some((1, 3, _))), Slice { .. } => false, _ => true, } From 6abf1a78a26172c9161e9427692f875ae01004ab Mon Sep 17 00:00:00 2001 From: Orson Peters Date: Fri, 20 Feb 2026 13:51:21 +0100 Subject: [PATCH 14/15] Add plan test --- py-polars/tests/unit/operations/test_top_k.py | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/py-polars/tests/unit/operations/test_top_k.py b/py-polars/tests/unit/operations/test_top_k.py index 6607c9185795..d70b8c619f74 100644 --- a/py-polars/tests/unit/operations/test_top_k.py +++ b/py-polars/tests/unit/operations/test_top_k.py @@ -1,3 +1,4 @@ +import re from collections.abc import Callable import pytest @@ -616,3 +617,12 @@ def test_top_k_union_null() -> None: pl.DataFrame({"a": [1, 2, 3, None, None]}, schema={"a": pl.Int64}), check_row_order=False, ) + + +def test_top_k_dyn_pred_pushdown() -> None: + df = pl.DataFrame({"x": [1], "y": [1]}) + plan = df.lazy().with_columns(pl.col.x * pl.col.x).sort("y").head(3).explain() + + pred = re.search(r"FILTER.*dynamic_predicate", plan) + with_cols = re.search(r"WITH_COLUMNS", plan) + assert pred.start() > with_cols.start() From c0e1329f5f2efb179578789e75933a66f68117dc Mon Sep 17 00:00:00 2001 From: Orson Peters Date: Fri, 20 Feb 2026 14:11:37 +0100 Subject: [PATCH 15/15] Mypy --- py-polars/tests/unit/operations/test_top_k.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/py-polars/tests/unit/operations/test_top_k.py b/py-polars/tests/unit/operations/test_top_k.py index d70b8c619f74..cab2be5ee132 100644 --- a/py-polars/tests/unit/operations/test_top_k.py +++ b/py-polars/tests/unit/operations/test_top_k.py @@ -625,4 +625,6 @@ def test_top_k_dyn_pred_pushdown() -> None: pred = re.search(r"FILTER.*dynamic_predicate", plan) with_cols = re.search(r"WITH_COLUMNS", plan) + assert pred is not None + assert with_cols is not None assert pred.start() > with_cols.start()