From 32dd30dd08bdb2d906659d3e3a900c2e97b71ed9 Mon Sep 17 00:00:00 2001 From: Qi Zhu <821684824@qq.com> Date: Fri, 9 Oct 2026 15:54:27 +0800 Subject: [PATCH] fix(physical-optimizer): decline the WindowTopN rewrite when the filter carries a fetch `WindowTopN` replaces the matched `FilterExec` outright, so anything the filter carries beyond its predicate has to be reproduced in the new plan or the rewrite has to decline. `FilterExec::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 a row limit over the filtered output. Today the fetch is dropped on the floor, so `FilterExec: rn <= 3, fetch=1` returns up to 3 rows per partition instead of 1 row overall. The default rule list runs `LimitPushdown` after `WindowTopN`, so no plan reaches the rule with a fetch set and this is not reachable from a stock pipeline. It is cheap insurance for pipelines that reorder the two, and the test pins the behaviour either way. --- .../tests/physical_optimizer/window_topn.rs | 30 +++++++++++++++++++ .../physical-optimizer/src/window_topn.rs | 15 ++++++++++ 2 files changed, 45 insertions(+) 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())?;