Repository navigation
Commit 176af6d
[SPARK-59995][4.3][SQL] Drop sort orders that hold a partition transform 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>1 parent c465689 commit 176af6d
7 files changed
Lines changed: 192 additions & 31 deletions
File tree
- sql
- catalyst/src/main
- java/org/apache/spark/sql/connector/read
- scala/org/apache/spark/sql/internal
- core/src
- main/scala/org/apache/spark/sql/execution/datasources/v2
- test/scala/org/apache/spark/sql
- connector
- execution/datasources/v2
Lines changed: 5 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
35 | 35 | | |
36 | 36 | | |
37 | 37 | | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
38 | 43 | | |
39 | 44 | | |
40 | 45 | | |
Lines changed: 6 additions & 4 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
2536 | 2536 | | |
2537 | 2537 | | |
2538 | 2538 | | |
2539 | | - | |
| 2539 | + | |
| 2540 | + | |
2540 | 2541 | | |
2541 | 2542 | | |
2542 | 2543 | | |
| |||
2548 | 2549 | | |
2549 | 2550 | | |
2550 | 2551 | | |
2551 | | - | |
2552 | | - | |
2553 | | - | |
| 2552 | + | |
| 2553 | + | |
| 2554 | + | |
| 2555 | + | |
2554 | 2556 | | |
2555 | 2557 | | |
2556 | 2558 | | |
| |||
Lines changed: 14 additions & 4 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
19 | 19 | | |
20 | 20 | | |
21 | 21 | | |
22 | | - | |
| 22 | + | |
23 | 23 | | |
24 | 24 | | |
25 | 25 | | |
| |||
116 | 116 | | |
117 | 117 | | |
118 | 118 | | |
| 119 | + | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
| 124 | + | |
| 125 | + | |
| 126 | + | |
119 | 127 | | |
120 | 128 | | |
| 129 | + | |
121 | 130 | | |
122 | 131 | | |
123 | | - | |
| 132 | + | |
| 133 | + | |
124 | 134 | | |
125 | 135 | | |
126 | | - | |
| 136 | + | |
127 | 137 | | |
128 | 138 | | |
129 | 139 | | |
130 | 140 | | |
131 | | - | |
| 141 | + | |
132 | 142 | | |
133 | 143 | | |
134 | 144 | | |
| |||
Lines changed: 14 additions & 17 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
24 | 24 | | |
25 | 25 | | |
26 | 26 | | |
27 | | - | |
28 | 27 | | |
29 | 28 | | |
30 | 29 | | |
| |||
361 | 360 | | |
362 | 361 | | |
363 | 362 | | |
364 | | - | |
365 | | - | |
| 363 | + | |
| 364 | + | |
366 | 365 | | |
367 | 366 | | |
368 | 367 | | |
| |||
374 | 373 | | |
375 | 374 | | |
376 | 375 | | |
377 | | - | |
| 376 | + | |
378 | 377 | | |
379 | 378 | | |
380 | 379 | | |
| |||
426 | 425 | | |
427 | 426 | | |
428 | 427 | | |
429 | | - | |
430 | | - | |
431 | | - | |
432 | | - | |
433 | | - | |
| 428 | + | |
| 429 | + | |
| 430 | + | |
434 | 431 | | |
435 | 432 | | |
436 | 433 | | |
| |||
510 | 507 | | |
511 | 508 | | |
512 | 509 | | |
513 | | - | |
514 | | - | |
515 | | - | |
516 | | - | |
| 510 | + | |
| 511 | + | |
| 512 | + | |
| 513 | + | |
517 | 514 | | |
518 | | - | |
| 515 | + | |
519 | 516 | | |
520 | 517 | | |
521 | | - | |
522 | | - | |
523 | | - | |
| 518 | + | |
| 519 | + | |
| 520 | + | |
524 | 521 | | |
Lines changed: 114 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
42 | 42 | | |
43 | 43 | | |
44 | 44 | | |
| 45 | + | |
45 | 46 | | |
46 | 47 | | |
47 | 48 | | |
| |||
5541 | 5542 | | |
5542 | 5543 | | |
5543 | 5544 | | |
| 5545 | + | |
| 5546 | + | |
| 5547 | + | |
| 5548 | + | |
| 5549 | + | |
| 5550 | + | |
| 5551 | + | |
| 5552 | + | |
| 5553 | + | |
| 5554 | + | |
| 5555 | + | |
| 5556 | + | |
| 5557 | + | |
| 5558 | + | |
| 5559 | + | |
| 5560 | + | |
| 5561 | + | |
| 5562 | + | |
| 5563 | + | |
| 5564 | + | |
| 5565 | + | |
| 5566 | + | |
| 5567 | + | |
| 5568 | + | |
| 5569 | + | |
| 5570 | + | |
| 5571 | + | |
| 5572 | + | |
| 5573 | + | |
| 5574 | + | |
| 5575 | + | |
| 5576 | + | |
| 5577 | + | |
| 5578 | + | |
| 5579 | + | |
| 5580 | + | |
| 5581 | + | |
| 5582 | + | |
| 5583 | + | |
| 5584 | + | |
| 5585 | + | |
| 5586 | + | |
| 5587 | + | |
| 5588 | + | |
| 5589 | + | |
| 5590 | + | |
| 5591 | + | |
| 5592 | + | |
| 5593 | + | |
| 5594 | + | |
| 5595 | + | |
| 5596 | + | |
| 5597 | + | |
| 5598 | + | |
| 5599 | + | |
| 5600 | + | |
| 5601 | + | |
| 5602 | + | |
| 5603 | + | |
| 5604 | + | |
| 5605 | + | |
| 5606 | + | |
| 5607 | + | |
| 5608 | + | |
| 5609 | + | |
| 5610 | + | |
| 5611 | + | |
| 5612 | + | |
| 5613 | + | |
| 5614 | + | |
| 5615 | + | |
| 5616 | + | |
| 5617 | + | |
| 5618 | + | |
| 5619 | + | |
| 5620 | + | |
| 5621 | + | |
| 5622 | + | |
| 5623 | + | |
| 5624 | + | |
| 5625 | + | |
| 5626 | + | |
| 5627 | + | |
| 5628 | + | |
| 5629 | + | |
| 5630 | + | |
| 5631 | + | |
| 5632 | + | |
| 5633 | + | |
| 5634 | + | |
| 5635 | + | |
| 5636 | + | |
| 5637 | + | |
| 5638 | + | |
| 5639 | + | |
| 5640 | + | |
| 5641 | + | |
| 5642 | + | |
| 5643 | + | |
| 5644 | + | |
| 5645 | + | |
| 5646 | + | |
| 5647 | + | |
| 5648 | + | |
| 5649 | + | |
| 5650 | + | |
| 5651 | + | |
| 5652 | + | |
| 5653 | + | |
| 5654 | + | |
| 5655 | + | |
| 5656 | + | |
| 5657 | + | |
5544 | 5658 | | |
5545 | 5659 | | |
5546 | 5660 | | |
| |||
Lines changed: 4 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1545 | 1545 | | |
1546 | 1546 | | |
1547 | 1547 | | |
1548 | | - | |
| 1548 | + | |
1549 | 1549 | | |
1550 | | - | |
| 1550 | + | |
1551 | 1551 | | |
1552 | 1552 | | |
| 1553 | + | |
| 1554 | + | |
1553 | 1555 | | |
1554 | 1556 | | |
0 commit comments