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
221 changes: 220 additions & 1 deletion datafusion/core/tests/physical_optimizer/window_topn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

//! Tests for the WindowTopN physical optimizer rule.

use std::collections::HashMap;
use std::sync::Arc;

use arrow::datatypes::{DataType, Field, Schema};
Expand All @@ -35,7 +36,7 @@ use datafusion_physical_optimizer::PhysicalOptimizerRule;
use datafusion_physical_optimizer::window_topn::WindowTopN;
use datafusion_physical_plan::collect;
use datafusion_physical_plan::displayable;
use datafusion_physical_plan::filter::FilterExec;
use datafusion_physical_plan::filter::{FilterExec, FilterExecBuilder};
use datafusion_physical_plan::metrics::MetricValue;
use datafusion_physical_plan::placeholder_row::PlaceholderRowExec;
use datafusion_physical_plan::projection::ProjectionExec;
Expand All @@ -52,6 +53,20 @@ fn schema() -> Arc<Schema> {
]))
}

/// Same shape as [`schema`], but every field carries metadata and the schema
/// itself carries top-level metadata. Used to check the rewrite reproduces the
/// filter's output schema exactly, not just its column names.
fn schema_with_metadata() -> Arc<Schema> {
let pk = Field::new("pk", DataType::Int64, false)
.with_metadata(HashMap::from([("k".to_string(), "pk".to_string())]));
let val = Field::new("val", DataType::Int64, false)
.with_metadata(HashMap::from([("k".to_string(), "val".to_string())]));
Arc::new(
Schema::new(vec![pk, val])
.with_metadata(HashMap::from([("origin".to_string(), "test".to_string())])),
)
}

fn plan_str(plan: &dyn ExecutionPlan) -> String {
displayable(plan).indent(true).to_string()
}
Expand Down Expand Up @@ -724,6 +739,210 @@ async fn partitioned_topk_exec_exposes_metrics() -> Result<()> {
.expect("output_batches metric")
.as_usize();
assert_eq!(output_batches, 5);
Ok(())
}

/// Regression: FilterExec that carries an embedded projection (e.g. from
/// an earlier filter/projection pushdown pass) used to make the rule bail
/// out entirely. The rule now captures the projection, applies the
/// PartitionedTopKExec rewrite, and wraps the result in a ProjectionExec
/// that reproduces the original output schema.
#[test]
fn filter_with_projection_still_rewrites() -> Result<()> {
let s = schema();
let input: Arc<dyn ExecutionPlan> = Arc::new(PlaceholderRowExec::new(Arc::clone(&s)));

let ordering = LexOrdering::new(vec![
PhysicalSortExpr::new_default(col("pk", &s)?).asc(),
PhysicalSortExpr::new_default(col("val", &s)?).asc(),
])
.unwrap();
let sort: Arc<dyn ExecutionPlan> =
Arc::new(SortExec::new(ordering, input).with_preserve_partitioning(true));

let partition_by = vec![col("pk", &s)?];
let order_by = vec![PhysicalSortExpr::new_default(col("val", &s)?).asc()];
let window_expr = Arc::new(StandardWindowExpr::new(
create_udwf_window_expr(
&row_number_udwf(),
&[],
&s,
"row_number".to_string(),
false,
)?,
&partition_by,
&order_by,
Arc::new(WindowFrame::new_bounds(
WindowFrameUnits::Rows,
WindowFrameBound::Preceding(ScalarValue::UInt64(None)),
WindowFrameBound::CurrentRow,
)),
));
let window: Arc<dyn ExecutionPlan> = Arc::new(BoundedWindowAggExec::try_new(
vec![window_expr],
sort,
InputOrderMode::Sorted,
true,
)?);

// Filter: row_number@2 <= 3, with an embedded projection that keeps
// only [pk, val] (drops the row_number column) — the shape produced
// by filter/projection pushdown when downstream doesn't need the
// window column.
let rn_col = Arc::new(Column::new("row_number", 2));
let limit_lit = lit(ScalarValue::UInt64(Some(3)));
let predicate = Arc::new(BinaryExpr::new(rn_col, Operator::LtEq, limit_lit));
let filter: Arc<dyn ExecutionPlan> = Arc::new(
FilterExecBuilder::new(predicate, window)
.apply_projection(Some(vec![0, 1]))?

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.

Optional: could we add a reordered or duplicate projection case such as [1, 0, 1] and assert that the optimized schema matches the original filter schema, including field and schema metadata? The current [0, 1] snapshot covers column removal, but this would strengthen coverage for ordering, duplicates, and metadata preservation.

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.

Thanks @kosiew — added in 7c3bc01 as reordered_duplicated_projection_preserves_filter_schema: a [1, 0, 1] projection (reorder + duplicate) over a schema with both field-level and schema-level metadata, asserting optimized.schema() == filter.schema().

It passes — the metadata does survive. Two guards so it can't pass for the wrong reason: it checks the rewrite actually fired (rather than the rule declining and handing back the input), and that FilterExec itself carries the metadata being compared (otherwise both sides would be empty).

.build()?,
);

let optimized = optimize(filter)?;
assert_snapshot!(plan_str(optimized.as_ref()), @r#"
ProjectionExec: expr=[pk@0 as pk, val@1 as val]
BoundedWindowAggExec: wdw=[row_number: Field { "row_number": UInt64 }, frame: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
PartitionedTopKExec: fn=row_number, fetch=3, partition=[pk@0], order=[val@1 ASC]
SortExec: expr=[pk@0 ASC, val@1 ASC], preserve_partitioning=[true]
PlaceholderRowExec
"#);
Ok(())
}

/// Builds the `FilterExec(row_number <= 3) → BoundedWindowAggExec → SortExec`
/// shape the rule targets, letting the caller decide what else the filter
/// carries.
fn window_topn_shape(
projection: Option<Vec<usize>>,
fetch: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
window_topn_shape_with_schema(&schema(), projection, fetch)
}

/// [`window_topn_shape`] over a caller-supplied schema.
fn window_topn_shape_with_schema(
s: &Arc<Schema>,
projection: Option<Vec<usize>>,
fetch: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
let s = Arc::clone(s);
let input: Arc<dyn ExecutionPlan> = Arc::new(PlaceholderRowExec::new(Arc::clone(&s)));
let ordering = LexOrdering::new(vec![
PhysicalSortExpr::new_default(col("pk", &s)?).asc(),
PhysicalSortExpr::new_default(col("val", &s)?).asc(),
])
.unwrap();
let sort: Arc<dyn ExecutionPlan> =
Arc::new(SortExec::new(ordering, input).with_preserve_partitioning(true));
let window_expr = Arc::new(StandardWindowExpr::new(
create_udwf_window_expr(
&row_number_udwf(),
&[],
&s,
"row_number".to_string(),
false,
)?,
&[col("pk", &s)?],
&[PhysicalSortExpr::new_default(col("val", &s)?).asc()],
Arc::new(WindowFrame::new_bounds(
WindowFrameUnits::Rows,
WindowFrameBound::Preceding(ScalarValue::UInt64(None)),
WindowFrameBound::CurrentRow,
)),
));
let window: Arc<dyn ExecutionPlan> = Arc::new(BoundedWindowAggExec::try_new(
vec![window_expr],
sort,
InputOrderMode::Sorted,
true,
)?);
let predicate = Arc::new(BinaryExpr::new(
Arc::new(Column::new("row_number", 2)),
Operator::LtEq,
lit(ScalarValue::UInt64(Some(3))),
));
let mut builder = FilterExecBuilder::new(predicate, window);
if let Some(indices) = projection {
builder = builder.apply_projection(Some(indices))?;
}
Ok(Arc::new(builder.with_fetch(fetch).build()?))
}

/// `FilterExec::fetch` is applied after the predicate, and the rewritten plan
/// has nowhere to carry it — a `PartitionedTopKExec` bounds rows *per
/// partition*, which is a different thing. The rule must decline rather than
/// return more rows than asked for.
#[test]
fn filter_with_projection_and_fetch_is_declined() -> Result<()> {
let filter = window_topn_shape(Some(vec![0, 1]), Some(1))?;
let optimized = optimize(Arc::clone(&filter))?;
assert_eq!(
plan_str(optimized.as_ref()),
plan_str(filter.as_ref()),
"a FilterExec carrying a fetch must be left alone"
);
Ok(())
}

/// The same guard covers the case that predates projection support: a filter
/// with `fetch` and no projection was previously rewritten with its fetch
/// silently dropped.
#[test]
fn filter_with_fetch_and_no_projection_is_declined() -> Result<()> {
let filter = window_topn_shape(None, Some(1))?;
let optimized = optimize(Arc::clone(&filter))?;
assert_eq!(
plan_str(optimized.as_ref()),
plan_str(filter.as_ref()),
"fetch must be preserved even without an embedded projection"
);
Ok(())
}

/// The rewrite reproduces `FilterExec`'s embedded projection as an outer
/// `ProjectionExec`, so the two must agree on more than column count: a
/// projection may reorder columns and repeat them, and the fields carry
/// metadata. `[1, 0, 1]` exercises all three at once.
#[test]
fn reordered_duplicated_projection_preserves_filter_schema() -> Result<()> {
let filter = window_topn_shape_with_schema(
&schema_with_metadata(),
Some(vec![1, 0, 1]),
None,
)?;
let filter_schema = filter.schema();
// Guard against a vacuous comparison: if `FilterExec` itself dropped the
// metadata, both sides would be empty and the assertion below would hold
// for the wrong reason.
assert_eq!(
filter_schema.metadata().get("origin").map(String::as_str),
Some("test"),
"filter must carry schema-level metadata for this test to mean anything"
);
assert_eq!(
filter_schema
.field(0)
.metadata()
.get("k")
.map(String::as_str),
Some("val"),
"filter must carry field metadata for this test to mean anything"
);

let optimized = optimize(Arc::clone(&filter))?;

// Guard against the assertion below passing vacuously because the rule
// declined the rewrite and handed back the input unchanged.
let optimized_str = plan_str(optimized.as_ref());
assert!(
optimized_str.starts_with("ProjectionExec:"),
"expected the rewrite to fire, got:\n{optimized_str}"
);
assert_eq!(
optimized.schema(),
filter_schema,
"rewritten plan must reproduce the filter's output schema exactly \
(order, duplicates, field metadata and schema metadata)"
);
Ok(())
}
55 changes: 50 additions & 5 deletions datafusion/physical-optimizer/src/window_topn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ use datafusion_common::config::ConfigOptions;
use datafusion_common::tree_node::{Transformed, TransformedResult, TreeNode};
use datafusion_common::{Result, ScalarValue};
use datafusion_expr::Operator;
use datafusion_physical_expr::PhysicalExpr;
use datafusion_physical_expr::expressions::{BinaryExpr, Column, Literal};
use datafusion_physical_expr::window::StandardWindowExpr;
use datafusion_physical_expr::{LexOrdering, PhysicalSortExpr};
Expand Down Expand Up @@ -142,11 +143,29 @@ impl WindowTopN {
// Step 1: Match FilterExec at the top
let filter = plan.downcast_ref::<FilterExec>()?;

// Don't handle filters with projections
if filter.projection().is_some() {
// The rewrite replaces the FilterExec entirely, so anything it carries
// beyond the predicate has to be reproduced or declined. `fetch` is
// applied by `FilterExec::execute` *after* the predicate, and the
// rewritten plan has nowhere to put it: a `PartitionedTopKExec` bounds
// rows per partition, which is not the same as a row limit over the
// filtered output. Decline rather than silently widen the result.
//
// Note this also guards a case that predates this rule's projection
// support: a FilterExec with `fetch` but no projection was already
// rewritten with its fetch dropped.
if filter.fetch().is_some() {
return None;
}

// A projection embedded in the FilterExec (from an earlier
// filter/projection pushdown pass) is captured here and re-applied
// via a wrapping ProjectionExec at the end so the rewrite preserves
// the original output schema.
let filter_projection: Option<Vec<usize>> = filter
.projection()
.as_ref()
.map(|p| p.iter().copied().collect());

// Step 2: Extract limit from predicate (rn <= K, rn < K, etc.)
let (col_idx, limit_n) = extract_window_limit(filter.predicate())?;

Expand Down Expand Up @@ -252,6 +271,34 @@ impl WindowTopN {
result = replace_children_if_necessary(node, vec![result]).ok()?;
}

// Step 10: Re-apply the FilterExec's embedded projection (if any)
// as an outer ProjectionExec. The projection indices refer to
// columns in `filter.input().schema()`, which equals `result`'s
// schema at this point (Steps 8-9 preserve schema), so the
// indices remain valid.
if let Some(indices) = filter_projection {

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.

Nice improvement capturing and restoring the embedded projection.

One thing I noticed is that the rewrite removes the FilterExec but does not preserve its fetch. FilterExecBuilder supports both an embedded projection and with_fetch, and FilterExec::execute applies the fetch after evaluating the predicate.

For example, if a matching projected filter has fetch=1, this rewrite currently produces only the outer ProjectionExec over the rewritten window. That returns all rn <= K rows instead of just one.

Could we either preserve the fetch with an equivalent outer limit/fetch operator, or skip this rewrite when filter.fetch().is_some()?

It would also be great to add a regression test covering the projection plus fetch case that executes the plan, or otherwise verifies that the row limit is preserved.

@zhuqi-lucas zhuqi-lucas Oct 4, 2026 •

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.

Good catch — fixed in f6bd9dd, and it turned out to be wider than the projection path.

I went with the second option: try_transform now declines when filter.fetch().is_some(). Preserving the fetch with an outer limit is not a straight substitution — PartitionedTopKExec bounds rows per partition, while FilterExec::fetch is a limit over the filtered output, so reproducing it would mean adding a real limit operator and reasoning about where it sits relative to the window. Declining keeps the rewrite honest and loses only the narrow intersection of "embedded projection and fetch".

Worth flagging on reachability: WindowTopN runs before LimitPushdown, so in the built-in pipeline the filter always has fetch: None — an outer LIMIT lands as its own GlobalLimitExec (new PROJ4 test). The guard is defensive rather than a live-bug fix. It isn't specific to the projection path either — try_transform on main never reads fetch, it just bails on projection().is_some() first.

Two regression tests, both asserting the plan comes back unchanged: filter_with_projection_and_fetch_is_declined and filter_with_fetch_and_no_projection_is_declined (the second is the pre-existing case).

Also rebased onto main and resolved the conflicts — find_window_below now returns a generic intermediates list rather than a single proj_between, so the rewrite composes with that. One existing snapshot needed updating: the SortExec under the window now stays in place, which the old expectation predated.

let input_schema = result.schema();
let field_count = input_schema.fields().len();
// Validate before indexing: an out-of-range index would panic
// in `input_schema.field(idx)`. Bail out of the rewrite instead
// so a malformed FilterExec projection can never crash the
// optimizer.
if indices.iter().any(|&idx| idx >= field_count) {
return None;
}
let projection_exprs: Vec<(Arc<dyn PhysicalExpr>, String)> = indices
.iter()
.map(|&idx| {
let field = input_schema.field(idx);
(
Arc::new(Column::new(field.name(), idx)) as Arc<dyn PhysicalExpr>,
field.name().clone(),
)
})
.collect();
result = Arc::new(ProjectionExec::try_new(projection_exprs, result).ok()?);
}

Some(result)
}
}
Expand Down Expand Up @@ -308,9 +355,7 @@ impl PhysicalOptimizerRule for WindowTopN {
/// - `10 >= rn` → `Some((2, 10))`
/// - `rn = 1` → `None` (equality not supported)
/// - `val <= 5` → `Some((1, 5))` (caller must verify it's a window column)
fn extract_window_limit(
predicate: &Arc<dyn datafusion_physical_expr::PhysicalExpr>,
) -> Option<(usize, usize)> {
fn extract_window_limit(predicate: &Arc<dyn PhysicalExpr>) -> Option<(usize, usize)> {
let binary = predicate.downcast_ref::<BinaryExpr>()?;
let op = binary.op();
let left = binary.left();
Expand Down
Loading
Loading