jayzhan211 commented on code in PR #25408:
URL: https://github.com/apache/datafusion/pull/25408#discussion_r4156310920
##########
datafusion/physical-optimizer/src/ensure_requirements/enforce_sorting/sort_pushdown.rs:
##########
@@ -566,8 +582,10 @@ fn pushdown_requirement_to_children(
} else {
parent_required.clone()
};
+
+ // Pushing a sort below an unordered fetch changes which rows it keeps.
Review Comment:
With the new limit branch, I can't find a built-in operator that reaches
this `return Ok(None)`. Limits are caught earlier,
`CoalescePartitionsExec`/`HashJoinExec` fail the
`maintains_input_order().all()` check, and SPM always has an ordering. When I
reverted this hunk, every unit test and the full sqllogictest suite still
passed. Either drop it, or add a unit test with a custom fetch-carrying
operator so the behavior is covered. A follow-up is fine.
##########
datafusion/physical-optimizer/src/ensure_requirements/enforce_sorting/sort_pushdown.rs:
##########
@@ -536,6 +536,22 @@ fn pushdown_requirement_to_children(
)
.then(|| Ok(vec![Some(parent_required)]))
.transpose()
+ } else if plan.is::<GlobalLimitExec>() || plan.is::<LocalLimitExec>() {
+ // Sorting below a limit changes which rows it keeps. Only refine a
+ // SortExec directly below it, which changes how ties are broken.
+ // Otherwise, sort the retained rows and let the input stop early.
+ let refines_sort_below = plan.children()[0].is::<SortExec>()
+ && plan.output_ordering().is_some_and(|ordering| {
+ plan.equivalence_properties().requirements_compatible(
+ parent_required.first().clone(),
+ ordering.clone().into(),
+ )
+ });
+ if refines_sort_below {
+ Ok(Some(vec![Some(parent_required)]))
+ } else {
+ Ok(None)
+ }
} else if plan.fetch().is_some()
Review Comment:
Not introduced here, but it's the same branch: an ordered
`SortPreservingMergeExec` with a fetch still gets the finer requirement pushed
below it, and its merge key stays `[b]`. The output is not ordered by `(b, a)`.
Main gives the same wrong result. Fine to handle in a follow-up issue.
```sql
set datafusion.execution.target_partitions = 4;
CREATE TABLE src (a INT, b INT) AS VALUES (1, 3), (2, 1), (3, 2), (4, 3),
(5, 1), (6, 2), (7, 3), (8, 1);
CREATE TABLE t AS SELECT a, b FROM src UNION ALL SELECT a + 10, b FROM src
UNION ALL SELECT a + 20, b FROM src UNION ALL SELECT a + 30, b FROM src;
SELECT * FROM (SELECT a, b FROM t ORDER BY b LIMIT 10) ORDER BY b, a;
-- 2 1 / 22 1 / 5 1 / 25 1 / ... (plan: SortPreservingMergeExec: [b],
fetch=10 -> SortExec: TopK [b, a])
```
Possible fix: decline the fetch branch for SPM (`&&
!is_sort_preserving_merge(plan)`), or rewrite the SPM's expressions when
pushing through it.
##########
datafusion/sqllogictest/test_files/limit.slt:
##########
@@ -362,6 +361,170 @@ SELECT * FROM (SELECT a FROM t1 LIMIT 4 OFFSET 3) ORDER
BY a;
6
7
+# LIMIT 4 OFFSET 3 keeps input rows (1, 7, 2, 10), then sorts them.
+statement ok
+CREATE TABLE unsorted (a INT) AS VALUES (5), (3), (9), (1), (7), (2), (10),
(4), (8), (6);
+
+query I
+SELECT * FROM (SELECT a FROM unsorted LIMIT 4 OFFSET 3) ORDER BY a;
+----
+1
+2
+7
+10
+
+# OFFSET skips the first 3 input rows (5, 3, 9) before sorting.
+query TT
+EXPLAIN SELECT * FROM (SELECT a FROM unsorted OFFSET 3) ORDER BY a;
+----
+logical_plan
+01)Sort: unsorted.a ASC NULLS LAST
+02)--Limit: skip=3, fetch=None
+03)----TableScan: unsorted projection=[a]
+physical_plan
+01)SortExec: expr=[a@0 ASC NULLS LAST], preserve_partitioning=[false]
+02)--GlobalLimitExec: skip=3, fetch=None
+03)----DataSourceExec: partitions=1, partition_sizes=[1]
+
+query I
+SELECT * FROM (SELECT a FROM unsorted OFFSET 3) ORDER BY a;
+----
+1
+2
+4
+6
+7
+8
+10
+
+# Even with input ordered on `a`, a TopK on `(a, b)` below the limit would
+# change which rows it keeps.
+query I
+COPY (SELECT column1 AS a, column2 AS b FROM (VALUES
+ (1, 9), (1, 8), (1, 7), (1, 6), (1, 5), (1, 4), (1, 3), (1, 2), (2, 1), (3,
0)))
+TO 'test_files/scratch/limit/ordered_prefix.parquet' STORED AS PARQUET;
+----
+10
+
+statement ok
+CREATE EXTERNAL TABLE ordered_prefix (a BIGINT, b BIGINT)
+STORED AS PARQUET
+WITH ORDER (a ASC)
+LOCATION 'test_files/scratch/limit/ordered_prefix.parquet';
+
+query TT
+EXPLAIN SELECT * FROM (SELECT a, b FROM ordered_prefix LIMIT 4 OFFSET 3) ORDER
BY a, b;
+----
+logical_plan
+01)Sort: ordered_prefix.a ASC NULLS LAST, ordered_prefix.b ASC NULLS LAST
+02)--Limit: skip=3, fetch=4
+03)----TableScan: ordered_prefix projection=[a, b], fetch=7
+physical_plan
+01)SortExec: expr=[a@0 ASC NULLS LAST, b@1 ASC NULLS LAST],
preserve_partitioning=[false]
+02)--GlobalLimitExec: skip=3, fetch=4
+03)----DataSourceExec: file_groups={1 group:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/limit/ordered_prefix.parquet]]},
projection=[a, b], limit=7, output_ordering=[a@0 ASC NULLS LAST],
file_type=parquet
+
+query II
+SELECT * FROM (SELECT a, b FROM ordered_prefix LIMIT 4 OFFSET 3) ORDER BY a, b;
+----
+1 3
+1 4
+1 5
+1 6
+
+# Refine the explicit sort to `(b, a)`. TopK keeps `skip + fetch` (3 + 4 = 7)
+# rows so the limit still returns 4.
+query TT
+EXPLAIN SELECT * FROM (SELECT a, b FROM ordered_prefix ORDER BY b LIMIT 4
OFFSET 3) ORDER BY b, a;
+----
+logical_plan
+01)Sort: ordered_prefix.b ASC NULLS LAST, ordered_prefix.a ASC NULLS LAST
+02)--Limit: skip=3, fetch=4
+03)----Sort: ordered_prefix.b ASC NULLS LAST, fetch=7
+04)------TableScan: ordered_prefix projection=[a, b]
+physical_plan
+01)GlobalLimitExec: skip=3, fetch=4
+02)--SortExec: TopK(fetch=7), expr=[b@1 ASC NULLS LAST, a@0 ASC NULLS LAST],
preserve_partitioning=[false]
+03)----DataSourceExec: file_groups={1 group:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/limit/ordered_prefix.parquet]]},
projection=[a, b], file_type=parquet, predicate=DynamicFilter [ empty ],
sort_order_for_reorder=[b@1 ASC NULLS LAST, a@0 ASC NULLS LAST],
dynamic_rg_pruning=eligible
+
+query II
+SELECT * FROM (SELECT a, b FROM ordered_prefix ORDER BY b LIMIT 4 OFFSET 3)
ORDER BY b, a;
+----
+1 3
+1 4
+1 5
+1 6
+
+# OFFSET without LIMIT: the refined sort is pushed into the explicit sort with
+# no fetch, and the limit still skips 3 rows.
+query TT
+EXPLAIN SELECT * FROM (SELECT a, b FROM ordered_prefix ORDER BY b OFFSET 3)
ORDER BY b, a;
+----
+logical_plan
+01)Sort: ordered_prefix.b ASC NULLS LAST, ordered_prefix.a ASC NULLS LAST
+02)--Limit: skip=3, fetch=None
+03)----Sort: ordered_prefix.b ASC NULLS LAST
+04)------TableScan: ordered_prefix projection=[a, b]
+physical_plan
+01)GlobalLimitExec: skip=3, fetch=None
+02)--SortExec: expr=[b@1 ASC NULLS LAST, a@0 ASC NULLS LAST],
preserve_partitioning=[false]
+03)----DataSourceExec: file_groups={1 group:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/limit/ordered_prefix.parquet]]},
projection=[a, b], file_type=parquet, sort_order_for_reorder=[b@1 ASC NULLS
LAST, a@0 ASC NULLS LAST]
+
+query II
+SELECT * FROM (SELECT a, b FROM ordered_prefix ORDER BY b OFFSET 3) ORDER BY
b, a;
+----
+1 3
+1 4
+1 5
+1 6
+1 7
+1 8
+1 9
+
+# A sort required by a window function, reaching the limit through a
projection,
+# also stays above LIMIT ... OFFSET so the input can stop after `skip + fetch`
+# rows. It used to be pushed below the limit as a TopK over the whole input
+# (#25378).
+statement ok
+CREATE TABLE window_over_limit (id INT, kind VARCHAR) AS VALUES
+ (19, 'a'), (18, 'a'), (17, 'a'), (16, 'a'), (15, 'a'),
+ (14, 'a'), (13, 'a'), (12, 'a'), (11, 'a'), (10, 'a'),
+ (9, 'a'), (8, 'a'), (7, 'a'), (6, 'a'), (5, 'a'),
+ (4, 'a'), (3, 'a'), (2, 'a'), (1, 'a'), (0, 'a');
+
+query TT
+EXPLAIN SELECT row_number() OVER (ORDER BY id) AS rn, id
+FROM (SELECT id FROM window_over_limit WHERE kind = 'a' LIMIT 5 OFFSET 5);
+----
+logical_plan
+01)Projection: row_number() ORDER BY [window_over_limit.id ASC NULLS LAST]
RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW AS rn, window_over_limit.id
+02)--WindowAggr: windowExpr=[[row_number() ORDER BY [window_over_limit.id ASC
NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW]]
+03)----Projection: window_over_limit.id
+04)------Limit: skip=5, fetch=5
+05)--------Filter: window_over_limit.kind = Utf8View("a")
+06)----------TableScan: window_over_limit projection=[id, kind]
+physical_plan
+01)ProjectionExec: expr=[row_number() ORDER BY [window_over_limit.id ASC NULLS
LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@1 as rn, id@0 as id]
+02)--BoundedWindowAggExec: wdw=[row_number() ORDER BY [window_over_limit.id
ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field {
"row_number() ORDER BY [window_over_limit.id ASC NULLS LAST] RANGE BETWEEN
UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED
PRECEDING AND CURRENT ROW], mode=[Sorted]
+03)----SortExec: expr=[id@0 ASC NULLS LAST], preserve_partitioning=[false]
+04)------ProjectionExec: expr=[id@0 as id]
+05)--------GlobalLimitExec: skip=5, fetch=5
+06)----------FilterExec: kind@1 = a, fetch=10
+07)------------DataSourceExec: partitions=1, partition_sizes=[1]
+
+query II
+SELECT row_number() OVER (ORDER BY id) AS rn, id
+FROM (SELECT id FROM window_over_limit WHERE kind = 'a' LIMIT 5 OFFSET 5);
+----
+1 10
+2 11
+3 12
+4 13
+5 14
+
+statement ok
+DROP TABLE window_over_limit;
+
Review Comment:
This PR also fixes a wrong-result bug on main: with multiple partitions, an
OFFSET-only subquery over `ORDER BY b` followed by `ORDER BY b, a` returns `1,
21, 4, 24, …` on main, because the `[b, a]` sort is pushed under
`SortPreservingMergeExec: [b]`. Please add a result check so this doesn't
regress:
```sql
statement ok
set datafusion.execution.target_partitions = 4;
statement ok
CREATE TABLE src (a INT, b INT) AS VALUES (1, 3), (2, 1), (3, 2), (4, 3),
(5, 1), (6, 2), (7, 3), (8, 1);
statement ok
CREATE TABLE t AS SELECT a, b FROM src UNION ALL SELECT a + 10, b FROM src
UNION ALL SELECT a + 20, b FROM src UNION ALL SELECT a + 30, b FROM src;
query II
SELECT * FROM (SELECT a, b FROM t ORDER BY b OFFSET 20) ORDER BY b, a;
----
1 3
4 3
7 3
11 3
14 3
17 3
21 3
24 3
27 3
31 3
34 3
37 3
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]