diff --git a/datafusion/physical-plan/src/lib.rs b/datafusion/physical-plan/src/lib.rs index 40489eada3e9d..77780c9f8e924 100644 --- a/datafusion/physical-plan/src/lib.rs +++ b/datafusion/physical-plan/src/lib.rs @@ -57,7 +57,7 @@ pub use crate::ordering::InputOrderMode; pub use crate::sort_pushdown::SortOrderPushdownResult; pub use crate::statistics::{ChildStats, StatisticsArgs, StatisticsContext}; pub use crate::stream::EmptyRecordBatchStream; -pub use crate::topk::TopK; +pub use crate::topk::{TopK, TopKDynamicFilters}; pub use crate::visitor::{ExecutionPlanVisitor, accept, visit_execution_plan}; pub use crate::work_table::WorkTable; pub use spill::spill_manager::SpillManager; diff --git a/datafusion/physical-plan/src/topk/mod.rs b/datafusion/physical-plan/src/topk/mod.rs index 05743472f9fe3..e5c2a40554e74 100644 --- a/datafusion/physical-plan/src/topk/mod.rs +++ b/datafusion/physical-plan/src/topk/mod.rs @@ -142,6 +142,47 @@ pub struct TopK { /// For more background, please also see the [Dynamic Filters: Passing Information Between Operators During Execution for 25x Faster Queries blog] /// /// [Dynamic Filters: Passing Information Between Operators During Execution for 25x Faster Queries blog]: https://datafusion.apache.org/blog/2025/09/10/dynamic-filters +/// +/// # Example +/// +/// Create a [`TopKDynamicFilters`] and pass it to [`TopK::try_new`]: +/// +/// ``` +/// # use std::sync::Arc; +/// # use arrow::datatypes::{DataType, Field, Schema}; +/// # use datafusion_execution::runtime_env::RuntimeEnv; +/// # use datafusion_physical_expr::{LexOrdering, PhysicalSortExpr}; +/// # use datafusion_physical_plan::expressions::{DynamicFilterPhysicalExpr, col, lit}; +/// # use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet; +/// # use parking_lot::RwLock; +/// use datafusion_physical_plan::{TopK, TopKDynamicFilters}; +/// +/// # fn main() -> datafusion_common::Result<()> { +/// let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])); +/// let sort_expr = PhysicalSortExpr::new_default(col("a", &schema)?); +/// +/// // The dynamic filter starts as `true` and is tightened as the TopK heap fills +/// let dynamic_filter = Arc::new(DynamicFilterPhysicalExpr::new( +/// vec![col("a", &schema)?], +/// lit(true), +/// )); +/// let filter = Arc::new(RwLock::new(TopKDynamicFilters::new(dynamic_filter))); +/// +/// let topk = TopK::try_new( +/// 0, // partition_id +/// Arc::clone(&schema), // schema +/// vec![], // common_sort_prefix +/// LexOrdering::from([sort_expr]), // expr +/// 10, // k +/// 8192, // batch_size +/// Arc::new(RuntimeEnv::default()), // runtime +/// &ExecutionPlanMetricsSet::new(), // metrics +/// filter, // filter +/// )?; +/// # let _ = topk; +/// # Ok(()) +/// # } +/// ``` #[derive(Debug)] pub struct TopKDynamicFilters { /// The current threshold shared by all TopK emitters that use this dynamic