diff --git a/datafusion/core/tests/physical_optimizer/window_topn.rs b/datafusion/core/tests/physical_optimizer/window_topn.rs index 85801e84c9189..5a04c31316de4 100644 --- a/datafusion/core/tests/physical_optimizer/window_topn.rs +++ b/datafusion/core/tests/physical_optimizer/window_topn.rs @@ -727,3 +727,33 @@ async fn partitioned_topk_exec_exposes_metrics() -> Result<()> { Ok(()) } + +/// `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 were asked for. +/// +/// Not reachable from the default rule list, which runs `LimitPushdown` after +/// `WindowTopN`; this pins the behaviour for pipelines that reorder the two. +#[test] +fn filter_with_fetch_is_declined() -> Result<()> { + let filter = build_window_topn_plan(3, Operator::LtEq)?; + + // Non-vacuity: without the fetch, this very shape is rewritten. + assert!( + find_partitioned_topk(&optimize(Arc::clone(&filter))?).is_some(), + "the fixture must be a shape the rule rewrites, or the assertion below \ + would hold for the wrong reason" + ); + + let with_fetch = filter + .with_fetch(Some(1)) + .expect("FilterExec accepts a fetch"); + let optimized = optimize(Arc::clone(&with_fetch))?; + assert_eq!( + plan_str(optimized.as_ref()), + plan_str(with_fetch.as_ref()), + "a FilterExec carrying a fetch must be left alone" + ); + Ok(()) +} diff --git a/datafusion/physical-optimizer/src/window_topn.rs b/datafusion/physical-optimizer/src/window_topn.rs index 758ad66de3017..9809b07ff2c3a 100644 --- a/datafusion/physical-optimizer/src/window_topn.rs +++ b/datafusion/physical-optimizer/src/window_topn.rs @@ -147,6 +147,21 @@ impl WindowTopN { return None; } + // 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. Declining keeps the + // rule from silently returning more rows than were asked for. + // + // The default rule list runs `LimitPushdown` after this rule, so no + // plan reaches here with a fetch today. The guard is cheap insurance + // for pipelines that reorder the two. + if filter.fetch().is_some() { + return None; + } + // Step 2: Extract limit from predicate (rn <= K, rn < K, etc.) let (col_idx, limit_n) = extract_window_limit(filter.predicate())?;