Skip to content

[SPARK-59995][SQL][4.2] Drop sort orders that hold a partition transform from a V2 scan's output ordering - #59304

Closed
peter-toth wants to merge 1 commit into
apache:branch-4.2from
peter-toth:SPARK-59995-transform-expression-codegen-4.2
Closed

peter-toth wants to merge 1 commit into
apache:branch-4.2from
peter-toth:SPARK-59995-transform-expression-codegen-4.2

Conversation

@peter-toth

Copy link
Copy Markdown
Contributor

Differences from #59251

This is the branch-4.2 backport of #59251. The description below is the original one.

  • The bug is on branch-4.2. 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.2.
  • docs/sql-performance-tuning.md is dropped. The config table on 4.2 has no rows for these configs.
  • BoundFunction.java is dropped. The Javadoc list it edits is not on 4.2.
  • 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.2.
  • In SQLConf.scala, the partitionKeyOrdering doc gets the new sentence on top of the 4.2 text. The 4.2 text has no clause about an ignored reported ordering.
  • In DataSourceV2ScanExecBase.scala, the code change is the same. The Scaladoc keeps the 4.2 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.2 does not have.
  • The derived k-way merge test uses V2_BUCKETING_ALLOW_JOIN_KEYS_SUBSET_OF_PARTITION_KEYS, the 4.2 name of V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS.
  • The derived k-way merge test also sets spark.sql.requireAllClusterKeysForCoPartition to false. On 4.2 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.2. So is spark.sql.execution.replaceHashWithSortAgg, which the scan ordering test turns on.
  • 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.2.
  • With the scan change kept and LazyRowOrdering built with GenerateOrdering again, the new GroupPartitionsExecSuite test fails.
  • Ran KeyGroupedPartitioningSuite (161), GroupPartitionsExecSuite (21), EnsureRequirementsSuite (58), WriteDistributionAndOrderingSuite (99), DataSourceV2Suite (54) and ProjectedOrderingAndPartitioningSuite (11). All 404 tests pass. TransformExpressionSuite is not on 4.2.
  • 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)

…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 dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+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.2 backport of #59251. The description below is the original one.

- The bug is on branch-4.2. 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.2.
- `docs/sql-performance-tuning.md` is dropped. The config table on 4.2 has no rows for these configs.
- `BoundFunction.java` is dropped. The Javadoc list it edits is not on 4.2.
- `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.2.
- In `SQLConf.scala`, the `partitionKeyOrdering` doc gets the new sentence on top of the 4.2 text. The 4.2 text has no clause about an ignored reported ordering.
- In `DataSourceV2ScanExecBase.scala`, the code change is the same. The Scaladoc keeps the 4.2 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.2 does not have.
- The derived k-way merge test uses `V2_BUCKETING_ALLOW_JOIN_KEYS_SUBSET_OF_PARTITION_KEYS`, the 4.2 name of `V2_BUCKETING_ALLOW_KEYS_SUBSET_OF_PARTITION_KEYS`.
- The derived k-way merge test also sets `spark.sql.requireAllClusterKeysForCoPartition` to `false`. On 4.2 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.2. So is `spark.sql.execution.replaceHashWithSortAgg`, which the scan ordering test turns on.
- 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.2.
- With the scan change kept and `LazyRowOrdering` built with `GenerateOrdering` again, the new `GroupPartitionsExecSuite` test fails.
- Ran `KeyGroupedPartitioningSuite` (161), `GroupPartitionsExecSuite` (21), `EnsureRequirementsSuite` (58), `WriteDistributionAndOrderingSuite` (99), `DataSourceV2Suite` (54) and `ProjectedOrderingAndPartitioningSuite` (11). All 404 tests pass. `TransformExpressionSuite` is not on 4.2.
- `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 #59304 from peter-toth/SPARK-59995-transform-expression-codegen-4.2.

Authored-by: Peter Toth <peter.toth@gmail.com>
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
@dongjoon-hyun

Copy link
Copy Markdown
Member

Merge Summary:

Posted by merge_spark_pr.py

@peter-toth

Copy link
Copy Markdown
Contributor Author

Thank you @dongjoon-hyun for the review!

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.

2 participants