Skip to content

Commit 6a3ebc3

Browse files
committed
test: add end-to-end SQL coverage for the embedded-projection rewrite
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.
1 parent e4ca432 commit 6a3ebc3

1 file changed

Lines changed: 112 additions & 0 deletions

File tree

‎datafusion/sqllogictest/test_files/window_topn.slt‎

Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1592,3 +1592,115 @@ set datafusion.execution.batch_size = 8192;
15921592

15931593
statement ok
15941594
set datafusion.optimizer.enable_window_topn = false;
1595+
1596+
###
1597+
### Embedded-projection group: the outer query drops `rn`, so projection
1598+
### pushdown folds the ProjectionExec into the FilterExec. Before this was
1599+
### supported the rule bailed on `filter.projection().is_some()` and the
1600+
### rewrite never fired for this (very common) shape.
1601+
###
1602+
1603+
statement ok
1604+
CREATE TABLE window_topn_t (id INT, pk INT, val INT) AS VALUES
1605+
(1, 1, 10),
1606+
(2, 1, 20),
1607+
(3, 1, 30),
1608+
(4, 1, 40),
1609+
(5, 2, 5),
1610+
(6, 2, 15),
1611+
(7, 2, 25),
1612+
(8, 3, 100),
1613+
(9, 3, 50),
1614+
(10, 3, 75);
1615+
1616+
statement ok
1617+
set datafusion.optimizer.enable_window_topn = true;
1618+
1619+
# Test PROJ1: results are correct when the filter carries a projection.
1620+
query III rowsort
1621+
SELECT id, pk, val FROM (
1622+
SELECT *, ROW_NUMBER() OVER (PARTITION BY pk ORDER BY val) as rn FROM window_topn_t
1623+
) WHERE rn <= 3;
1624+
----
1625+
1 1 10
1626+
10 3 75
1627+
2 1 20
1628+
3 1 30
1629+
5 2 5
1630+
6 2 15
1631+
7 2 25
1632+
8 3 100
1633+
9 3 50
1634+
1635+
statement ok
1636+
SET datafusion.explain.physical_plan_only = true;
1637+
1638+
# Test PROJ2: the plan is rewritten to PartitionedTopKExec, and the filter's
1639+
# projection is preserved as the top-level ProjectionExec.
1640+
query TT
1641+
EXPLAIN SELECT id, pk, val FROM (
1642+
SELECT *, ROW_NUMBER() OVER (PARTITION BY pk ORDER BY val) as rn FROM window_topn_t
1643+
) WHERE rn <= 3;
1644+
----
1645+
physical_plan
1646+
01)ProjectionExec: expr=[id@0 as id, pk@1 as pk, val@2 as val]
1647+
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [window_topn_t.pk] ORDER BY [window_topn_t.val ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [window_topn_t.pk] ORDER BY [window_topn_t.val ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
1648+
03)----RepartitionExec: partitioning=Hash([pk@1], 4), input_partitions=1, maintains_sort_order=true
1649+
04)------PartitionedTopKExec: fn=row_number, fetch=3, partition=[pk@1], order=[val@2 ASC NULLS LAST]
1650+
05)--------DataSourceExec: partitions=1, partition_sizes=[1]
1651+
1652+
# Test PROJ3: with the rule off, the same query shows the shape the rule now
1653+
# handles -- a FilterExec carrying `projection=[...]` over the window.
1654+
statement ok
1655+
set datafusion.optimizer.enable_window_topn = false;
1656+
1657+
query TT
1658+
EXPLAIN SELECT id, pk, val FROM (
1659+
SELECT *, ROW_NUMBER() OVER (PARTITION BY pk ORDER BY val) as rn FROM window_topn_t
1660+
) WHERE rn <= 3;
1661+
----
1662+
physical_plan
1663+
01)FilterExec: row_number() PARTITION BY [window_topn_t.pk] ORDER BY [window_topn_t.val ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 <= 3, projection=[id@0, pk@1, val@2]
1664+
02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [window_topn_t.pk] ORDER BY [window_topn_t.val ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [window_topn_t.pk] ORDER BY [window_topn_t.val ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
1665+
03)----SortExec: expr=[pk@1 ASC NULLS LAST, val@2 ASC NULLS LAST], preserve_partitioning=[false]
1666+
04)------DataSourceExec: partitions=1, partition_sizes=[1]
1667+
1668+
statement ok
1669+
set datafusion.optimizer.enable_window_topn = true;
1670+
1671+
# Test PROJ4: adding an outer LIMIT must not lose the limit. `WindowTopN` runs
1672+
# before `LimitPushdown`, so the rewrite happens while the filter still has no
1673+
# fetch, and the limit lands as its own `GlobalLimitExec` -- it is never folded
1674+
# into a filter this rule then discards. `target_partitions = 1` is what makes
1675+
# `LimitPushdown` able to reach the filter at all.
1676+
statement ok
1677+
SET datafusion.execution.target_partitions = 1;
1678+
1679+
query TT
1680+
EXPLAIN SELECT id, pk, val FROM (
1681+
SELECT *, ROW_NUMBER() OVER (PARTITION BY pk ORDER BY val) as rn FROM window_topn_t
1682+
) WHERE rn <= 3 LIMIT 2;
1683+
----
1684+
physical_plan
1685+
01)ProjectionExec: expr=[id@0 as id, pk@1 as pk, val@2 as val]
1686+
02)--GlobalLimitExec: skip=0, fetch=2
1687+
03)----BoundedWindowAggExec: wdw=[row_number() PARTITION BY [window_topn_t.pk] ORDER BY [window_topn_t.val ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [window_topn_t.pk] ORDER BY [window_topn_t.val ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
1688+
04)------PartitionedTopKExec: fn=row_number, fetch=3, partition=[pk@1], order=[val@2 ASC NULLS LAST]
1689+
05)--------DataSourceExec: partitions=1, partition_sizes=[1]
1690+
1691+
statement ok
1692+
SET datafusion.explain.physical_plan_only = false;
1693+
1694+
query III
1695+
SELECT id, pk, val FROM (
1696+
SELECT *, ROW_NUMBER() OVER (PARTITION BY pk ORDER BY val) as rn FROM window_topn_t
1697+
) WHERE rn <= 3 LIMIT 2;
1698+
----
1699+
1 1 10
1700+
2 1 20
1701+
1702+
statement ok
1703+
SET datafusion.execution.target_partitions = 4;
1704+
1705+
statement ok
1706+
set datafusion.optimizer.enable_window_topn = false;

0 commit comments

Comments
 (0)