Skip to content

[SPARK-58411][SQL] Preserve ORDER BY in LIMIT decorrelation when correlation references only the outer table - #57614

Open
AveryQi115 wants to merge 2 commits into
apache:masterfrom
AveryQi115:SPARK-58411
Open

[SPARK-58411][SQL] Preserve ORDER BY in LIMIT decorrelation when correlation references only the outer table#57614
AveryQi115 wants to merge 2 commits into
apache:masterfrom
AveryQi115:SPARK-58411

Conversation

@AveryQi115

@AveryQi115 AveryQi115 commented Jul 28, 2026

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

In DecorrelateInnerQuery, the Limit case peels the Sort off the subquery and saves its ordering. When the correlated predicate references only outer columns, decorrelation produces no partitioning key (partitionFields.isEmpty) and that branch re-emits the limit as Limit(k, newChild) / Limit(k, Offset(m, newChild)) without re-applying the saved ordering, turning ORDER BY ... LIMIT k into an arbitrary LIMIT k.

This PR re-applies the saved ordering as a global Sort below the limit in the partitionFields.isEmpty branch (covering both LIMIT and LIMIT ... OFFSET), so ordered limits stay order-preserving and deterministic. A new internal flag spark.sql.optimizer.decorrelateLimitOffsetLegacyIncorrectOrderHandling.enabled (default false) can restore the legacy behavior.

The standalone Offset case (no LIMIT) is not affected: it decorrelates the original input with the Sort retained, so its ordering is already preserved.

Why are the changes needed?

The dropped ORDER BY makes an ordered-limit subquery return a non-deterministic row. For example:

SELECT t1a, t1b FROM t1
WHERE t1c = (SELECT t2c FROM t2 WHERE t1b < t1d ORDER BY t2c LIMIT 1);

With distinct t2c values {12, 16, NULL} and Spark's default ASC NULLS FIRST, the correct ordered LIMIT 1 returns NULL, so the row is correctly excluded. With the bug, an arbitrary row is picked, which can incorrectly emit a row and varies by physical execution order.

ground truth verification

image Result from Postgresql is [val1a | 16]. This is because the default ordering in Postgresql is asc with NULL last. But in Spark sql, it is asc with null first.

Does this PR introduce any user-facing change?

Yes. Correlated scalar/predicate subqueries of the form ORDER BY ... LIMIT whose correlation references only the outer table now return correct, deterministic results. Results that previously depended on the dropped ordering may change. The legacy behavior can be restored via the config above.

How was this patch tested?

New unit tests in DecorrelateInnerQuerySuite covering default ordering, explicit ASC NULLS LAST, LIMIT ... OFFSET, and the legacy-flag path. Golden-file results for scalar-subquery-predicate.sql are updated accordingly.

Was this patch authored or co-authored using generative AI tooling?

Yes, drafted with assistance from Claude.

AveryQi115 and others added 2 commits July 28, 2026 21:19
…elation references only the outer table

### What changes were proposed in this pull request?

In `DecorrelateInnerQuery`, the `Limit` case peels the `Sort` off the subquery
and saves its ordering. When the correlated predicate references only outer
columns, decorrelation produces no partitioning key (`partitionFields.isEmpty`)
and that branch re-emits the limit as `Limit(k, newChild)` /
`Limit(k, Offset(m, newChild))` without re-applying the saved ordering, turning
`ORDER BY ... LIMIT k` into an arbitrary `LIMIT k`.

This PR re-applies the saved ordering as a global `Sort` below the limit in the
`partitionFields.isEmpty` branch (covering both `LIMIT` and `LIMIT ... OFFSET`),
so ordered limits stay order-preserving and deterministic. A new internal flag
`spark.sql.optimizer.decorrelateLimitOffsetLegacyIncorrectOrderHandling.enabled`
(default false) can restore the legacy behavior.

The standalone `Offset` case (no `LIMIT`) is not affected: it decorrelates the
original input with the `Sort` retained, so its ordering is already preserved.

### Why are the changes needed?

The dropped `ORDER BY` makes an ordered-limit subquery return a non-deterministic
row. For example:

```
SELECT t1a, t1b FROM t1
WHERE t1c = (SELECT t2c FROM t2 WHERE t1b < t1d ORDER BY t2c LIMIT 1);
```

With distinct `t2c` values `{12, 16, NULL}` and Spark's default `ASC NULLS FIRST`,
the correct ordered `LIMIT 1` returns `NULL`, so the row is correctly excluded.
With the bug, an arbitrary row is picked, which can incorrectly emit a row and
varies by physical execution order.

### Does this PR introduce any user-facing change?

Yes. Correlated scalar/predicate subqueries of the form `ORDER BY ... LIMIT`
whose correlation references only the outer table now return correct,
deterministic results. Results that previously depended on the dropped ordering
may change. The legacy behavior can be restored via the config above.

### How was this patch tested?

New unit tests in `DecorrelateInnerQuerySuite` covering default ordering,
explicit `ASC NULLS LAST`, `LIMIT ... OFFSET`, and the legacy-flag path.
Golden-file results for `scalar-subquery-predicate.sql` are updated accordingly.

### Was this patch authored or co-authored using generative AI tooling?

Yes, drafted with assistance from Claude.
@AveryQi115

Copy link
Copy Markdown
Contributor Author

@agubichev @srielau @cloud-fan for review.

Plz let me know for such correctness issues, do we add flags to guard it? Do we need to backport the fix? Would like some suggestions for such fix since I'm not familiar with the process. TY

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant