Skip to content

feat(physical-optimizer): support FilterExec with embedded projection in WindowTopN - #23599

Open
zhuqi-lucas wants to merge 6 commits into
apache:mainfrom
zhuqi-lucas:qizhu/window-topn-filter-projection
Open

zhuqi-lucas wants to merge 6 commits into
apache:mainfrom
zhuqi-lucas:qizhu/window-topn-filter-projection

Conversation

@zhuqi-lucas

@zhuqi-lucas zhuqi-lucas commented Jul 15, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Rationale for this change

WindowTopN::try_transform currently bails out unconditionally when the top FilterExec carries an embedded projection (filter.projection().is_some()):

// Don't handle filters with projections
if filter.projection().is_some() {
    return None;
}

In practice, DataFusion's filter/projection pushdown pass often collapses a downstream ProjectionExec INTO the FilterExec's projection field, so real plans that match the FilterExec → BoundedWindowAggExec → SortExec pattern in every other respect are silently skipped and fall back to a full sort.

Observed at our prod:

SELECT ..., ROW_NUMBER() OVER (PARTITION BY g, d ORDER BY t DESC) AS rn
FROM ...
WHERE rn = 1

matches this exact pattern, but the FilterExec ends up with a projection pushed in, so WindowTopN skips and the sort-based path runs. This defeats the point of #21479 for a common shape.

What changes are included in this PR?

  • Capture the FilterExec's projection indices at the start of try_transform.
  • Run the existing PartitionedTopKExec rewrite as usual.
  • Re-apply the captured projection as an outer ProjectionExec at the end so the transformed plan preserves the original output schema.

Plan shape before (for a filter with projection = [0, 1] that drops the ROW_NUMBER column):

FilterExec(predicate=rn <= 3, projection=[0, 1])
  BoundedWindowAggExec(ROW_NUMBER PBY pk OBY val)
    SortExec(pk, val)

Plan shape after:

ProjectionExec(expr=[pk@0, val@1])
  BoundedWindowAggExec(ROW_NUMBER PBY pk OBY val)
    PartitionedTopKExec(fetch=3, partition=[pk], order=[val], fn=row_number)

Are these changes tested?

Yes. Added a regression test filter_with_projection_still_rewrites in datafusion/core/tests/physical_optimizer/window_topn.rs:

  • Builds FilterExec(rn <= 3, projection=[0, 1]) → BoundedWindowAggExec → SortExec.
  • Runs WindowTopN with enable_window_topn = true.
  • Asserts (via insta::assert_snapshot!) the resulting plan is ProjectionExec → BoundedWindowAggExec → PartitionedTopKExec → PlaceholderRowExec.

Existing 13 tests continue to pass:

cargo test -p datafusion --test core_integration physical_optimizer::window_topn -- --nocapture
test result: ok. 14 passed; 0 failed; 0 ignored; 0 measured

Also verified cargo fmt --all and cargo build -p datafusion-physical-optimizer succeed.

Are there any user-facing changes?

Yes. Queries matching ROW_NUMBER() OVER (PARTITION BY ...) ... WHERE rn OP K where an earlier optimizer pass has embedded a projection into the FilterExec will now be rewritten to PartitionedTopKExec (previously they silently fell back to a full sort). Only enabled when datafusion.optimizer.enable_window_topn = true.

Copilot AI review requested due to automatic review settings July 15, 2026 06:38
@github-actions github-actions Bot added optimizer Optimizer rules core Core DataFusion crate labels Jul 15, 2026

Copilot AI left a comment

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.

Pull request overview

Enhances the WindowTopN physical optimizer rule to rewrite eligible plans even when FilterExec contains an embedded projection (introduced by earlier filter/projection pushdown), preserving the original output schema by re-applying that projection as an outer ProjectionExec.

Changes:

  • Capture FilterExec embedded projection indices during WindowTopN::try_transform and re-apply them as an outer ProjectionExec after the rewrite.
  • Generalize extract_window_limit’s predicate parameter type by importing PhysicalExpr.
  • Add a regression test covering FilterExec(predicate, projection=[..]) → BoundedWindowAggExec → SortExec rewrites.

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

File Description
datafusion/physical-optimizer/src/window_topn.rs Preserve FilterExec embedded projections by wrapping the rewritten plan in a ProjectionExec.
datafusion/core/tests/physical_optimizer/window_topn.rs Add regression test ensuring embedded projections don’t prevent the PartitionedTopKExec rewrite.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment thread datafusion/physical-optimizer/src/window_topn.rs

@kosiew kosiew left a comment •

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.

@zhuqi-lucas,

Thanks for working on this.
Supporting FilterExec with embedded projections is a nice improvement, and the regression test covers the plan shape well.

I found one correctness issue that should be addressed before merging. The rewrite currently drops FilterExec::fetch, which can change query results.

// 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.

@github-actions

github-actions Bot commented Oct 4, 2026

Copy link
Copy Markdown

Thank you for your contribution. Unfortunately, this pull request is stale because it has been open 60 days with no activity. Please remove the stale label or comment or this will be closed in 7 days.

@github-actions github-actions Bot added the Stale PR has not had any activity for some time label Oct 4, 2026
@zhuqi-lucas
zhuqi-lucas force-pushed the qizhu/window-topn-filter-projection branch from 2eb9b29 to d75a081 Compare October 4, 2026 14:44
@zhuqi-lucas

Copy link
Copy Markdown
Contributor Author

@zhuqi-lucas,

Thanks for working on this. Supporting FilterExec with embedded projections is a nice improvement, and the regression test covers the plan shape well.

I found one correctness issue that should be addressed before merging. The rewrite currently drops FilterExec::fetch, which can change query results.

Thank you @kosiew for review, addressed review comments and added more tests now.

@zhuqi-lucas
zhuqi-lucas requested a review from kosiew October 4, 2026 14:46
@codecov-commenter

codecov-commenter commented Oct 4, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 91.66667% with 2 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.74%. Comparing base (3ed377a) to head (130b549).
⚠️ Report is 3 commits behind head on main.

Files with missing lines Patch % Lines
datafusion/physical-optimizer/src/window_topn.rs 91.66% 1 Missing and 1 partial ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #23599      +/-   ##
==========================================
- Coverage   82.74%   82.74%   -0.01%     
==========================================
  Files        1147     1147              
  Lines      449767   449787      +20     
  Branches   449767   449787      +20     
==========================================
+ Hits       372160   372168       +8     
- Misses      54938    54946       +8     
- Partials    22669    22673       +4     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@github-actions github-actions Bot removed the Stale PR has not had any activity for some time label Oct 5, 2026

@kosiew kosiew left a comment

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.

@zhuqi-lucas,

Thanks for working on this. I revalidated the cumulative changes and did not find any blocking issues. The earlier FilterExec fetch-loss concern is addressed, and the projection rewrite looks sound.

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).

@2010YOUY01

Copy link
Copy Markdown
Contributor

Thank you for the great explanation and the fix!

Is there any SQL test we can add here, for example an EXPLAIN ... test showing that the plan can now be optimized to WindowTopK, whereas previously it could not because of the embedded projection?

I suspect this may not be reproducible with the default optimizer rule order, since WindowTopN runs before ProjectionPushdown. Only in the latter rewrite can a projection be fused into a sibling node, in this case FilterExec. Your downstream optimizer rule list seems to include some deeper customization that makes this case possible.

This should not be a blocker, though. It seems more like a hidden assumption than a specified rule. I think this points to a deeper problem that makes optimizer rules harder to maintain in general, and I'm thinking about how we could address it systematically.

@zhuqi-lucas
zhuqi-lucas force-pushed the qizhu/window-topn-filter-projection branch from d75a081 to 6a3ebc3 Compare October 6, 2026 13:26
@github-actions github-actions Bot added the sqllogictest SQL Logic Tests (.slt) label Oct 6, 2026
@zhuqi-lucas

zhuqi-lucas commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor Author

Thanks @2010YOUY01 — added a PROJ group in window_topn.slt.

Superseded: these tests pass on main unchanged, so they demonstrated nothing and have been removed. See #23599 (comment) — ProjectionPushdown runs after WindowTopN, so the embedded projection they showed was created once the rule had already declined. Replaced with an e2e test that fails on main.

@2010YOUY01

Copy link
Copy Markdown
Contributor

I tried to add this new slt to main, and they're passing 🤔 If it's regression test, it's expected to fail (not optimized to windowTopK) without this PR.

This claim should be right #23599 (comment), there is a implicit 'no-embedded-projection' zone in the default optimizer rule list. And WindowTopN is assuming no projection in filter currently, to simplify its implementation.

        let rules: Vec<Arc<dyn PhysicalOptimizerRule + Send + Sync>> = vec![
            // ---- BEGIN: no-embedded-projection zone ----
            Arc::new(OutputRequirements::new_add_mode()),
            // ......
            Arc::new(WindowTopN::new()),
            // ......
            // ---- END: no-embedded-projection zone ----
            Arc::new(ProjectionPushdown::new()),
        ]

This hidden convention itself is not a blocker of this PR, but the real issue is, fix and improvements should verifiable with e2e tests, and this optimizer improvement can only be verified with UT taking a plan shape, that can't get produced from the default optimizer rule order.

It feel a bit like DataFusion core being extended to adapt downstream custom optimizer pipeline, I'm uncertain if it's expected. Would love to hear your thoughts on this.

@zhuqi-lucas

Copy link
Copy Markdown
Contributor Author

Thanks @2010YOUY01 — you're right, and I confirmed it: I dropped the new slt onto main unchanged and the whole file passes. It isn't a regression test, so I've removed it. The trap was that EXPLAIN shows the final plan, not what WindowTopN saw — ProjectionPushdown embeds the projection after the rule has already run.

One more data point for your "hidden convention" framing: no other rule in physical-optimizer reads filter.projection() at all. So this PR wouldn't just rely on the convention — it'd be the first rule in core to step outside it.

On your actual question. DataFusion is a library for building engines, so I don't think "a downstream pipeline produces this shape" is automatically out of scope. But the bar I'd want is that core could plausibly produce the shape itself, and today it can't — the only thing keeping it from happening is an ordering nobody wrote down. So I'd rather make that explicit than quietly depend on it.

Concretely, I'd split this:

  • The fetch guard is independent of all this and holds on its own — happy to keep it in a separate PR.
  • For the projection support, I'll follow your call. If it's useful to core, I have an e2e test that does discriminate (it fails on main and passes here) by inserting an extra ProjectionPushdown before WindowTopN and planning real SQL — though that still needs a non-default order, which is exactly your point. If you'd rather not take it, I'm fine carrying it downstream.

Either way the convention deserves to be written down — a rule silently ceasing to fire has no failing test anywhere. Happy to help with whatever systematic fix you have in mind.

… in WindowTopN

The WindowTopN rule bailed out unconditionally when the FilterExec at
the top of the pattern carried an embedded projection (from an earlier
filter/projection pushdown pass), even when the underlying pattern
otherwise matched. That skipped rewrite path caused ROW_NUMBER
top-K-per-group queries to fall back to a full sort in production.

Capture the FilterExec's projection indices at the start, run the
existing PartitionedTopKExec rewrite as usual, and re-apply the
captured projection as an outer ProjectionExec so the transformed plan
preserves the original output schema.

Adds a regression test in datafusion/core/tests/physical_optimizer/window_topn.rs
covering the FilterExec-with-projection shape.
FilterExec::execute applies fetch after the predicate, and the rewritten plan
has nowhere to carry it: PartitionedTopKExec bounds rows per partition, which
is not a row limit over the filtered output. Dropping it silently returned
more rows than asked for.

The guard also covers a case that predates this PR: a filter with fetch and no
embedded projection was already rewritten with its fetch dropped.
main now rebuilds intermediate nodes generically, so the SortExec under the
window stays in place; the old snapshot predated that.
Adds a PROJ group to window_topn.slt covering the shape the rule now
handles: the outer query drops `rn`, so projection pushdown folds the
ProjectionExec into the FilterExec.

- PROJ1 asserts results are correct through the rewrite.
- PROJ2 shows the plan is rewritten to PartitionedTopKExec.
- PROJ3 shows the same query with the rule off, where the FilterExec
  still carries `projection=[...]` -- the shape that used to make the
  rule bail.
- PROJ4 covers an outer LIMIT, confirming the limit survives as its own
  GlobalLimitExec rather than being folded into a discarded filter.
…, duplicated projection

The rewrite reproduces FilterExec's embedded projection as an outer
ProjectionExec, so the two must agree on more than column count. Adds a
`[1, 0, 1]` case -- reorder, duplicate and metadata in one -- over a
schema carrying both field-level and schema-level metadata, and asserts
the optimized schema equals the filter's.

The test guards against two ways of passing for the wrong reason: it
checks the rewrite actually fired, and that FilterExec itself carries
the metadata being compared.
…on main

The slt added earlier passes on main unchanged, so it demonstrated
nothing: EXPLAIN shows the final plan, not what WindowTopN saw.
ProjectionPushdown runs after WindowTopN in the default list, so the
embedded projection those tests showed was created once the rule had
already declined. Removed.

Replaced with an e2e test over public API that does discriminate: it
inserts an extra ProjectionPushdown pass directly before WindowTopN,
leaving the rest of the default list intact, then plans real SQL. The
filter then reaches the rule carrying projection=[0, 1], which main
declines (falling back to a full sort) and this branch rewrites.

The test asserts up front that the default list still orders the two
rules the other way round, so it reports that it has become meaningless
rather than silently passing if that ever changes.
@zhuqi-lucas
zhuqi-lucas force-pushed the qizhu/window-topn-filter-projection branch from 7c3bc01 to 130b549 Compare October 8, 2026 13:28
@github-actions github-actions Bot removed the sqllogictest SQL Logic Tests (.slt) label Oct 8, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core Core DataFusion crate optimizer Optimizer rules

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants