Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion datafusion/physical-plan/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add a compiling rustdoc example or external integration test that imports both types from the crate root, constructs TopKDynamicFilters, and passes it to TopK::try_new? Tests inside the private topk module would still pass without this re-export, so a public-path compilation test would catch this visibility issue if it regresses.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, this is a good idea. I'll try to get around to this later today or tomorrow

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I had claude add a trivial doc example and validated that reverting my fix causes the doc test to fail. Let me know what you think, I'm happy to make changes.

pub use crate::visitor::{ExecutionPlanVisitor, accept, visit_execution_plan};
pub use crate::work_table::WorkTable;
pub use spill::spill_manager::SpillManager;
Expand Down
41 changes: 41 additions & 0 deletions datafusion/physical-plan/src/topk/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading