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
30 changes: 30 additions & 0 deletions datafusion/core/tests/physical_optimizer/window_topn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}
15 changes: 15 additions & 0 deletions datafusion/physical-optimizer/src/window_topn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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())?;

Expand Down
Loading