Repository navigation
[SPARK-59995][SQL][4.3] Drop sort orders that hold a partition transform from a V2 scan's output ordering - #59303
Closed
peter-toth wants to merge 1 commit into
Conversation
…rom a V2 scan's output ordering
### What changes were proposed in this pull request?
`DataSourceV2ScanExecBase.outputOrdering` now drops every sort order that holds a partition transform (`TransformExpression`):
- in a reported ordering, a transform ends the leading run of sort orders the scan keeps, like a sort order over a pruned column;
- a transform on a partition key is dropped too, while a sort order on another partition key still holds after it, since each partition holds a single key;
- in an ordering the scan derives from its partition keys, the transform keys are left out.
Dropping them loses nothing today when Spark can call the transform's function, since no operator then requires an ordering over the transform. The write path sorts by the function call instead (`DistributionAndOrderingUtils`). Dropping them is also a safe way to handle a transform Spark cannot evaluate, and it saves comparisons nobody uses. If an ordering over a transform becomes a real requirement, this needs to be revisited.
The k-way merge of `GroupPartitionsExec` now also builds its comparator with `RowOrdering.create`, like `SortExec`. It still generates code, and falls back to interpreted evaluation when that fails. It is still built on the executor, at the first comparison. `LazyCodeGenOrdering` is renamed `LazyRowOrdering` to match.
The docs of `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled`, its 4.4 migration note and the `SupportsReportOrdering` Javadoc now say that partition transforms are left out. The `BoundFunction` Javadoc now names the scan merge, instead of a retained reported ordering, as what a semantic `equals` is needed for.
### Why are the changes needed?
With `spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled` on, `GroupPartitionsExec` can coalesce partitions with a k-way merge over its child's ordering (SPARK-55715). The merge compares rows with generated code (`LazyCodeGenOrdering`). A `TransformExpression` generates no code. So when the ordering contained a partition transform, such as `years(arrive_time)`, the task failed:
```
[INTERNAL_ERROR] Cannot generate code for expression: transformexpression(org.apache.spark.sql.connector.catalog.functions.YearsFunction$..., input[2, timestamp, true], None) SQLSTATE: XX000
```
The ordering can contain a transform in two ways:
- the source reports it via `SupportsReportOrdering`, e.g. `[id, name, years(arrive_time)]`;
- the scan derives it from its partition keys (SPARK-56241), e.g. `[id, years(arrive_time)]`.
### Does this PR introduce _any_ user-facing change?
Yes.
- A query that k-way merges over a transform used to fail with the error above. Now it returns its rows. The merge still runs when the sort orders the scan keeps satisfy the parent.
- A sort order on a partition key after a transform can now lead the scan's ordering. For example, with partition keys `[years(ts), id]` the scan reports `[id]` instead of `[years(ts), id]`. So a sort on `id` right above the scan, e.g. from `sortWithinPartitions("id")`, is no longer needed. The new ordering can also turn a hash aggregate into a sort aggregate, since `spark.sql.execution.replaceHashWithSortAgg` keys off the scan's ordering. For example, `SELECT id, max(ts) FROM t GROUP BY id` over the partition keys `[years(ts), id]` now plans its partial aggregate as a `SortAggregateExec`.
- The failure exists since 4.2.0, which has both the k-way merge (SPARK-55715) and the key-derived ordering (SPARK-56241). It needs `spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled`, which is off by default on every branch. The derived case also needs `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled`, which defaults to true only on master and branch-4.x.
### How was this patch tested?
- New tests in `KeyGroupedPartitioningSuite`:
- "a scan's output ordering drops sort orders that hold a partition transform" checks the three cases above, and that a sort on the kept key needs no `SortExec`;
- two k-way merge tests, one for each way above. Each checks the answer, that the plan k-way merges, and the ordering it merges over. They fail on master with the error above. The derived one uses `days`, a function Spark cannot call.
- New test in `GroupPartitionsExecSuite`: "the k-way merge ordering falls back to interpreted evaluation". It serializes a `LazyRowOrdering` over a transform, which generates no code, and compares rows with it. It fails without the fallback.
- The SPARK-56321 test in `WriteDistributionAndOrderingSuite` now checks the reported ordering, and that the output ordering drops its transform.
### Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 5.5)
Closes apache#59251 from peter-toth/SPARK-59995-transform-expression-codegen.
Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Peter Toth <peter.toth@gmail.com>
(cherry picked from commit 1af0063)
dongjoon-hyun
approved these changes
Oct 9, 2026
dongjoon-hyun
left a comment
Member
There was a problem hiding this comment.
+1, LGTM. Thank you, @peter-toth .
dongjoon-hyun
pushed a commit
that referenced
this pull request
Oct 9, 2026
…orm from a V2 scan's output ordering ### Differences from #59251 This is the branch-4.3 backport of #59251. The description below is the original one. - The bug is on branch-4.3. Without the fix, both new k-way merge tests fail with the `Cannot generate code for expression` error. - `docs/sql-migration-guide.md` is dropped. Its hunks edit the 4.4 notes on the `partitionKeyOrdering` and `preserveKeyOrderingOnCoalesce` default changes (SPARK-59396), which are not on 4.3. - `docs/sql-performance-tuning.md` is dropped. The config table on 4.3 has no rows for these configs. - `BoundFunction.java` is dropped. The Javadoc list it edits is not on 4.3. - `SupportsReportOrdering.java` only gets the new paragraph about transforms. The paragraph before it on master comes from a later change that is not on 4.3. - In `SQLConf.scala`, the `partitionKeyOrdering` doc gets the new sentence on top of the 4.3 text. The 4.3 text has no clause about an ignored reported ordering. - In `DataSourceV2ScanExecBase.scala`, the code change is the same. The Scaladoc keeps the 4.3 first paragraph and adds the new one. - `GroupPartitionsExec.scala` and `WriteDistributionAndOrderingSuite.scala` apply cleanly. - In `GroupPartitionsExecSuite.scala`, only the imports conflicted. The new test is the same. - `KeyGroupedPartitioningSuite.scala` adds the `SortAggregateExec` import, which 4.3 does not have. - The derived k-way merge test also sets `spark.sql.requireAllClusterKeysForCoPartition` to `false`. On 4.3 a join on a subset of the partition keys needs it. Without it the plan does not k-way merge. - The other new tests already set the configs they need. `partitionKeyOrdering`, `preserveKeyOrderingOnCoalesce` and `preserveOrderingOnCoalesce` are all off by default on 4.3. - Without the fix, the 3 new `KeyGroupedPartitioningSuite` tests and the adjusted SPARK-56321 test fail. For this run `DataSourceV2ScanExecBase.scala` and `GroupPartitionsExec.scala` were restored from branch-4.3. - With the scan change kept and `LazyRowOrdering` built with `GenerateOrdering` again, the new `GroupPartitionsExecSuite` test fails. - Ran `KeyGroupedPartitioningSuite` (188), `GroupPartitionsExecSuite` (23), `EnsureRequirementsSuite` (58), `WriteDistributionAndOrderingSuite` (99), `DataSourceV2Suite` (69) and `ProjectedOrderingAndPartitioningSuite` (36). All 473 tests pass. `TransformExpressionSuite` is not on 4.3. - `dev/lint-scala` passes. ### What changes were proposed in this pull request? `DataSourceV2ScanExecBase.outputOrdering` now drops every sort order that holds a partition transform (`TransformExpression`): - in a reported ordering, a transform ends the leading run of sort orders the scan keeps, like a sort order over a pruned column; - a transform on a partition key is dropped too, while a sort order on another partition key still holds after it, since each partition holds a single key; - in an ordering the scan derives from its partition keys, the transform keys are left out. Dropping them loses nothing today when Spark can call the transform's function, since no operator then requires an ordering over the transform. The write path sorts by the function call instead (`DistributionAndOrderingUtils`). Dropping them is also a safe way to handle a transform Spark cannot evaluate, and it saves comparisons nobody uses. If an ordering over a transform becomes a real requirement, this needs to be revisited. The k-way merge of `GroupPartitionsExec` now also builds its comparator with `RowOrdering.create`, like `SortExec`. It still generates code, and falls back to interpreted evaluation when that fails. It is still built on the executor, at the first comparison. `LazyCodeGenOrdering` is renamed `LazyRowOrdering` to match. The docs of `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled`, its 4.4 migration note and the `SupportsReportOrdering` Javadoc now say that partition transforms are left out. The `BoundFunction` Javadoc now names the scan merge, instead of a retained reported ordering, as what a semantic `equals` is needed for. ### Why are the changes needed? With `spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled` on, `GroupPartitionsExec` can coalesce partitions with a k-way merge over its child's ordering (SPARK-55715). The merge compares rows with generated code (`LazyCodeGenOrdering`). A `TransformExpression` generates no code. So when the ordering contained a partition transform, such as `years(arrive_time)`, the task failed: ``` [INTERNAL_ERROR] Cannot generate code for expression: transformexpression(org.apache.spark.sql.connector.catalog.functions.YearsFunction$..., input[2, timestamp, true], None) SQLSTATE: XX000 ``` The ordering can contain a transform in two ways: - the source reports it via `SupportsReportOrdering`, e.g. `[id, name, years(arrive_time)]`; - the scan derives it from its partition keys (SPARK-56241), e.g. `[id, years(arrive_time)]`. ### Does this PR introduce _any_ user-facing change? Yes. - A query that k-way merges over a transform used to fail with the error above. Now it returns its rows. The merge still runs when the sort orders the scan keeps satisfy the parent. - A sort order on a partition key after a transform can now lead the scan's ordering. For example, with partition keys `[years(ts), id]` the scan reports `[id]` instead of `[years(ts), id]`. So a sort on `id` right above the scan, e.g. from `sortWithinPartitions("id")`, is no longer needed. The new ordering can also turn a hash aggregate into a sort aggregate, since `spark.sql.execution.replaceHashWithSortAgg` keys off the scan's ordering. For example, `SELECT id, max(ts) FROM t GROUP BY id` over the partition keys `[years(ts), id]` now plans its partial aggregate as a `SortAggregateExec`. - The failure exists since 4.2.0, which has both the k-way merge (SPARK-55715) and the key-derived ordering (SPARK-56241). It needs `spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled`, which is off by default on every branch. The derived case also needs `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled`, which defaults to true only on master and branch-4.x. ### How was this patch tested? - New tests in `KeyGroupedPartitioningSuite`: - "a scan's output ordering drops sort orders that hold a partition transform" checks the three cases above, and that a sort on the kept key needs no `SortExec`; - two k-way merge tests, one for each way above. Each checks the answer, that the plan k-way merges, and the ordering it merges over. They fail on master with the error above. The derived one uses `days`, a function Spark cannot call. - New test in `GroupPartitionsExecSuite`: "the k-way merge ordering falls back to interpreted evaluation". It serializes a `LazyRowOrdering` over a transform, which generates no code, and compares rows with it. It fails without the fallback. - The SPARK-56321 test in `WriteDistributionAndOrderingSuite` now checks the reported ordering, and that the output ordering drops its transform. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5.5) Closes #59303 from peter-toth/SPARK-59995-transform-expression-codegen-4.3. Authored-by: Peter Toth <peter.toth@gmail.com> Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
Member
|
Merge Summary:
Posted by |
Contributor
Author
|
Thank you @dongjoon-hyun. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Differences from #59251
This is the branch-4.3 backport of #59251. The description below is the original one.
Cannot generate code for expressionerror.docs/sql-migration-guide.mdis dropped. Its hunks edit the 4.4 notes on thepartitionKeyOrderingandpreserveKeyOrderingOnCoalescedefault changes (SPARK-59396), which are not on 4.3.docs/sql-performance-tuning.mdis dropped. The config table on 4.3 has no rows for these configs.BoundFunction.javais dropped. The Javadoc list it edits is not on 4.3.SupportsReportOrdering.javaonly gets the new paragraph about transforms. The paragraph before it on master comes from a later change that is not on 4.3.SQLConf.scala, thepartitionKeyOrderingdoc gets the new sentence on top of the 4.3 text. The 4.3 text has no clause about an ignored reported ordering.DataSourceV2ScanExecBase.scala, the code change is the same. The Scaladoc keeps the 4.3 first paragraph and adds the new one.GroupPartitionsExec.scalaandWriteDistributionAndOrderingSuite.scalaapply cleanly.GroupPartitionsExecSuite.scala, only the imports conflicted. The new test is the same.KeyGroupedPartitioningSuite.scalaadds theSortAggregateExecimport, which 4.3 does not have.spark.sql.requireAllClusterKeysForCoPartitiontofalse. On 4.3 a join on a subset of the partition keys needs it. Without it the plan does not k-way merge.partitionKeyOrdering,preserveKeyOrderingOnCoalesceandpreserveOrderingOnCoalesceare all off by default on 4.3.KeyGroupedPartitioningSuitetests and the adjusted SPARK-56321 test fail. For this runDataSourceV2ScanExecBase.scalaandGroupPartitionsExec.scalawere restored from branch-4.3.LazyRowOrderingbuilt withGenerateOrderingagain, the newGroupPartitionsExecSuitetest fails.KeyGroupedPartitioningSuite(188),GroupPartitionsExecSuite(23),EnsureRequirementsSuite(58),WriteDistributionAndOrderingSuite(99),DataSourceV2Suite(69) andProjectedOrderingAndPartitioningSuite(36). All 473 tests pass.TransformExpressionSuiteis not on 4.3.dev/lint-scalapasses.What changes were proposed in this pull request?
DataSourceV2ScanExecBase.outputOrderingnow drops every sort order that holds a partition transform (TransformExpression):Dropping them loses nothing today when Spark can call the transform's function, since no operator then requires an ordering over the transform. The write path sorts by the function call instead (
DistributionAndOrderingUtils). Dropping them is also a safe way to handle a transform Spark cannot evaluate, and it saves comparisons nobody uses. If an ordering over a transform becomes a real requirement, this needs to be revisited.The k-way merge of
GroupPartitionsExecnow also builds its comparator withRowOrdering.create, likeSortExec. It still generates code, and falls back to interpreted evaluation when that fails. It is still built on the executor, at the first comparison.LazyCodeGenOrderingis renamedLazyRowOrderingto match.The docs of
spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled, its 4.4 migration note and theSupportsReportOrderingJavadoc now say that partition transforms are left out. TheBoundFunctionJavadoc now names the scan merge, instead of a retained reported ordering, as what a semanticequalsis needed for.Why are the changes needed?
With
spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabledon,GroupPartitionsExeccan coalesce partitions with a k-way merge over its child's ordering (SPARK-55715). The merge compares rows with generated code (LazyCodeGenOrdering). ATransformExpressiongenerates no code. So when the ordering contained a partition transform, such asyears(arrive_time), the task failed:The ordering can contain a transform in two ways:
SupportsReportOrdering, e.g.[id, name, years(arrive_time)];[id, years(arrive_time)].Does this PR introduce any user-facing change?
Yes.
[years(ts), id]the scan reports[id]instead of[years(ts), id]. So a sort onidright above the scan, e.g. fromsortWithinPartitions("id"), is no longer needed. The new ordering can also turn a hash aggregate into a sort aggregate, sincespark.sql.execution.replaceHashWithSortAggkeys off the scan's ordering. For example,SELECT id, max(ts) FROM t GROUP BY idover the partition keys[years(ts), id]now plans its partial aggregate as aSortAggregateExec.spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled, which is off by default on every branch. The derived case also needsspark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled, which defaults to true only on master and branch-4.x.How was this patch tested?
KeyGroupedPartitioningSuite:SortExec;days, a function Spark cannot call.GroupPartitionsExecSuite: "the k-way merge ordering falls back to interpreted evaluation". It serializes aLazyRowOrderingover a transform, which generates no code, and compares rows with it. It fails without the fallback.WriteDistributionAndOrderingSuitenow checks the reported ordering, and that the output ordering drops its transform.Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 5.5)