From a265ecd9fb4d2b7f3b9fd20bb43410d20f77ff6e Mon Sep 17 00:00:00 2001 From: Mason Hall Date: Thu, 17 Sep 2026 15:23:31 -0400 Subject: [PATCH 1/2] Make `TopKDynamicFilters` public This is needed as an argument to `TopK::try_new()`, so without it it's impossible to construct a `TopK`. --- datafusion/physical-plan/src/lib.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datafusion/physical-plan/src/lib.rs b/datafusion/physical-plan/src/lib.rs index 5ff6cec374ee1..47c57b265bcc8 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; From f2b09228e5cb22a02fb67ea3d7f8b2364811535c Mon Sep 17 00:00:00 2001 From: Mason Hall Date: Wed, 7 Oct 2026 16:19:12 -0400 Subject: [PATCH 2/2] add a doc example to ensure that `TopKDynamicFilters` stays publicly accessible --- datafusion/physical-plan/src/topk/mod.rs | 41 ++++++++++++++++++++++++ 1 file changed, 41 insertions(+) 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