Repository navigation
feat(physical-optimizer): support FilterExec with embedded projection in WindowTopN - #23599
zhuqi-lucas wants to merge 5 commits into
Conversation
There was a problem hiding this comment.
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
FilterExecembedded projection indices duringWindowTopN::try_transformand re-apply them as an outerProjectionExecafter the rewrite. - Generalize
extract_window_limit’s predicate parameter type by importingPhysicalExpr. - Add a regression test covering
FilterExec(predicate, projection=[..]) → BoundedWindowAggExec → SortExecrewrites.
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.
9148eb2 to
2eb9b29
Compare
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
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. |
2eb9b29 to
d75a081
Compare
Thank you @kosiew for review, addressed review comments and added more tests now. |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #23599 +/- ##
==========================================
- Coverage 82.72% 82.72% -0.01%
==========================================
Files 1147 1147
Lines 448093 448113 +20
Branches 448093 448113 +20
==========================================
+ Hits 370693 370698 +5
- Misses 54891 54900 +9
- Partials 22509 22515 +6 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
kosiew
left a comment
There was a problem hiding this comment.
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]))? |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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).
|
Thank you for the great explanation and the fix! Is there any SQL test we can add here, for example an I suspect this may not be reproducible with the default optimizer rule order, since 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. |
… 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.
d75a081 to
6a3ebc3
Compare
|
Thanks @2010YOUY01 — added in 6a3ebc3, a
|
…, 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.
Which issue does this PR close?
Rationale for this change
WindowTopN::try_transformcurrently bails out unconditionally when the topFilterExeccarries an embedded projection (filter.projection().is_some()):In practice, DataFusion's filter/projection pushdown pass often collapses a downstream
ProjectionExecINTO the FilterExec'sprojectionfield, so real plans that match theFilterExec → BoundedWindowAggExec → SortExecpattern in every other respect are silently skipped and fall back to a full sort.Observed at our prod:
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?
try_transform.ProjectionExecat 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):Plan shape after:
Are these changes tested?
Yes. Added a regression test
filter_with_projection_still_rewritesindatafusion/core/tests/physical_optimizer/window_topn.rs:FilterExec(rn <= 3, projection=[0, 1])→BoundedWindowAggExec→SortExec.WindowTopNwithenable_window_topn = true.insta::assert_snapshot!) the resulting plan isProjectionExec → BoundedWindowAggExec → PartitionedTopKExec → PlaceholderRowExec.Existing 13 tests continue to pass:
Also verified
cargo fmt --allandcargo build -p datafusion-physical-optimizersucceed.Are there any user-facing changes?
Yes. Queries matching
ROW_NUMBER() OVER (PARTITION BY ...) ... WHERE rn OP Kwhere an earlier optimizer pass has embedded a projection into the FilterExec will now be rewritten toPartitionedTopKExec(previously they silently fell back to a full sort). Only enabled whendatafusion.optimizer.enable_window_topn = true.