Repository navigation
[SPARK-59995][SQL] Drop sort orders that hold a partition transform from a V2 scan's output ordering - #59251
[SPARK-59995][SQL] Drop sort orders that hold a partition transform from a V2 scan's output ordering#59251peter-toth wants to merge 5 commits into
Conversation
… bound function `TransformExpression.doGenCode` now generates the code of the bound function call that `eval` runs. It still throws when there is no such call. The k-way merge of `GroupPartitionsExec` compares rows with generated code. An ordering with a partition transform, reported by the source or derived from the partition keys, made the task fail with `Cannot generate code for expression`.
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for the fix, @peter-toth. I reviewed the latest commit, 791e616, and left 12 inline comments. Here is a summary, ordered by priority.
Correctness
TransformExpressionL194: the k-way merge is still planned when Spark cannot build the call (e.g. this suite's owndays), so the task fails where aSortExecon the join keys would succeed.TransformExpressionL192: theCodegenFallbackthatresolvedFunctionmay hold is invisible toCollapseCodegenStages, which theCodegenFallback.generatescaladoc asks such a node to be named in.TransformExpressionL193:LazyCodeGenOrderingstill has no interpreted fallback, so a call thatevalruns but the generated comparator cannot compile still fails the task.TransformExpressionL195:resolvedFunctionskips the casts thatBoundFunction.inputTypes()promises; this predates the PR, and the generated code now inherits it.
Performance
KeyGroupedPartitioningSuiteL6598: no requirement can contain a transform, so the merge evaluates the transform key only as a tie-breaker nobody consumes; cuttingkWayMergeOrderingbefore it would also avoid item 1.
Tests
KeyGroupedPartitioningSuiteL6591: neither test checks that the merge orders rows by the transform's value, and the first one never runsyears().KeyGroupedPartitioningSuiteL6630: no end-to-end test covers aproduceResult-only transform.TransformExpressionSuiteL130: no null input, although the fixtures differ for null.TransformExpressionSuiteL129: nothing checks thatdoGenCodeemits the call's own code rather than anevalfallback.
Docs and comments
TransformExpressionL189: "onlyproduceResult" is narrower than when the call falls back, and theresolvedFunctionscaladoc does not cover the new ordering reader.
Cleanup
TransformExpressionSuiteL148: the new fixtures do not overridecanonicalName(), whose default is a random UUID.TransformExpressionSuiteL137: the reduced-keys test can useainstead of a secondBoundReference.
Items 1 and 2 matter most, and item 5 may be the simplest way to address item 1.
| case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this) | ||
| } | ||
|
|
||
| // `resolvedFunction` is not a child, so a tree walk does not see it. For a function with only |
There was a problem hiding this comment.
Nit: "a function with only produceResult" is narrower than when this falls back. V2ExpressionUtils.findMethod uses getDeclaredMethod, so the call is also an ApplyFunctionExpression when invoke is inherited or declared with other parameter classes; the PR description's "declares ... in its own class" is the accurate wording. Separately, the resolvedFunction scaladoc (L169-176) explains why its reduced-keys arm is unreachable only for partitionings and the write path ("every consumer of a reduced partitioning refuses it first"), while doGenCode now reaches it through an ordering. It is unreachable there because only GroupPartitionsExec's output partitioning reports a reduced transform, and the only ordering built from partition keys is the scan's own.
Could we say "when Spark calls the function through produceResult" here, and add the ordering case to that scaladoc?
There was a problem hiding this comment.
Thanks. The doGenCode change is gone in 913f317. No output ordering holds a transform now, so the resolvedFunction scaladoc stays as it is.
| // `resolvedFunction` is not a child, so a tree walk does not see it. For a function with only | ||
| // `produceResult` it falls back to `eval` on the input row, which whole-stage codegen does not | ||
| // provide (see `CodegenFallback.generate`). No whole-stage operator generates code for a | ||
| // transform today. |
There was a problem hiding this comment.
resolvedFunction is not a child, so the CodegenFallback it may hold is invisible to CollapseCodegenStages. For a function called through produceResult, fn is an ApplyFunctionExpression, whose generated code is ((Expression) references[i]).eval(ctx.INPUT_ROW). The CodegenFallback.generate scaladoc (CodegenFallback.scala:47-50) says that a node falling back for a reason other than holding such a descendant "would have to be named in CollapseCodegenStages as well, since nothing else would tell it that the input row this code needs is not there", but CollapseCodegenStages.supportCodegen (WholeStageCodegenExec.scala:918-927) only walks children. I agree that no CodegenSupport operator holds a transform today, but the first one that does would fail worse than before: consume sets ctx.INPUT_ROW = null (WholeStageCodegenExec.scala:188), so ProjectExec/FilterExec would fail code generation with an NPE from the code interpolator, and where INPUT_ROW names another row (HashAggregateExec's unsafeRowBuffer, a join's build row) the call would read the wrong row silently. Before this PR the same plan failed at codegen time with INTERNAL_ERROR.
Could we add case _: TransformExpression => false to CollapseCodegenStages.supportCodegen (it changes nothing today) and mention TransformExpression in the CodegenFallback.generate scaladoc, instead of relying on this comment?
There was a problem hiding this comment.
The doGenCode change is reverted in 913f317. A transform in whole-stage codegen fails with INTERNAL_ERROR again, as on master. So no CollapseCodegenStages change is needed.
| // `produceResult` it falls back to `eval` on the input row, which whole-stage codegen does not | ||
| // provide (see `CodegenFallback.generate`). No whole-stage operator generates code for a | ||
| // transform today. | ||
| override protected def doGenCode(ctx: CodegenContext, ev: ExprCode): ExprCode = |
There was a problem hiding this comment.
The merge comparator still has no interpreted fallback, so this fixes one expression rather than the mechanism that failed. LazyCodeGenOrdering (GroupPartitionsExec.scala:806-810) calls GenerateOrdering.generate directly, while SortExec.createSorter uses RowOrdering.create, a CodeGeneratorWithInterpretedFallback that falls back to InterpretedOrdering on any NonFatal codegen error under the default spark.sql.codegen.factoryMode. So a transform whose call eval can run, but whose generated code cannot compile, still fails the task. For example, take a Java connector function declared as a package-private class with a public magic invoke: Invoke.eval works, since commons-lang3 MethodUtils.getMatchingAccessibleMethod makes that method accessible, but the comparator casts to the class, which Janino should reject as inaccessible from the generated class's package. The same function works in a plain query or a SortExec, which log a warning and fall back. I understand codegen was chosen because the comparator is on the hot path (#55116), and the sql/core tests run with CODEGEN_ONLY (SparkSessionBinder.scala:85), so they would keep exercising the generated path.
Could LazyCodeGenOrdering use RowOrdering.create(sortOrders, schema), or catch NonFatal itself, as a defense in depth?
There was a problem hiding this comment.
Taken in 913f317. LazyCodeGenOrdering is now LazyRowOrdering, and it builds its comparator with RowOrdering.create. That still happens in a @transient lazy val, so the comparator is still built on the executor. A new test in GroupPartitionsExecSuite serializes a LazyRowOrdering over a transform and checks that it compares rows under FALLBACK.
| // transform today. | ||
| override protected def doGenCode(ctx: CodegenContext, ev: ExprCode): ExprCode = | ||
| throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this) | ||
| resolvedFunction match { |
There was a problem hiding this comment.
The k-way merge is still planned when Spark cannot build this call, so such a query still fails at task time. GroupPartitionsExec.kWayMergeIsFeasible (GroupPartitionsExec.scala:248-249) only checks child.outputOrdering.nonEmpty && childIsSafeForKWayMerge, and nothing at planning time checks that the ordering's transforms can be evaluated. This suite's own days is such a function: DaysFunction declares neither invoke nor produceResult in its own class, and V2ExpressionUtils.findMethod uses getDeclaredMethod, so building resolvedFunction here throws SCALAR_FUNCTION_NOT_FULLY_IMPLEMENTED. If the second new test partitions by days("arrive_time") instead of years("arrive_time"), the merge is planned over [id, days(arrive_time)] and the first comparison on the executor throws, while without the merge EnsureRequirements would add a SortExec on id and the query would return its rows. A bound function that is not a ScalarFunction fails the same way, with INTERNAL_ERROR from L196. This is not a regression, and the PR description scopes the fix to "as long as Spark can build the call", but ScalarFunction's Javadoc says the magic method is resolved during query analysis, while here it is first resolved on the executors.
Could kWayMergeIsFeasible refuse a merge whose ordering holds a transform without a buildable call, resolving it on the driver, and could we add a days test that checks the plan falls back to a SortExec and returns the right rows?
There was a problem hiding this comment.
Fixed in 913f317. The merge no longer compares rows by a transform, so it never calls the function. The derived-ordering test now uses days. The plan keeps the merge, over [id].
| override protected def doGenCode(ctx: CodegenContext, ev: ExprCode): ExprCode = | ||
| throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this) | ||
| resolvedFunction match { | ||
| case Some(fn) => fn.genCode(ctx) |
There was a problem hiding this comment.
Minor: this runs the call without the casts that BoundFunction.inputTypes() promises, now in generated code too. resolvedFunction (L177-182) passes the transform's raw children to V2ExpressionUtils.resolveScalarFunction, and nothing analyzes the resulting Invoke/StaticInvoke/ApplyFunctionExpression, so the cast promised in BoundFunction's Javadoc ("Spark will cast input values to the required data types") never happens; only the write path coerces them (TypeCoercionExecutor in DistributionAndOrderingUtils.scala:86). With the test catalog, bucket(4, c) over an INT column fails with a ClassCastException in ApplyFunctionExpression's reusedRow, whose slot is typed by the declared LongType, and years(d) over a DATE column, which UnboundYearsFunction accepts, reads the day count as microseconds. In the merge this only affects a tie-breaker that no parent relies on, but the keyed shuffle (ShuffleExchangeExec.scala:487) evaluates the same call. This predates the PR, and the generated code behaves exactly like eval.
Could resolvedFunction cast its arguments to inputTypes() as the write path does, or would you prefer a separate JIRA for it?
There was a problem hiding this comment.
Agreed, this predates the PR. After 913f317 the merge no longer evaluates a transform, but the keyed shuffle still does. I filed SPARK-60027 for it.
|
|
||
| test("SPARK-59995: a transform whose keys a join reduced generates no code") { | ||
| // The call no longer computes such keys, so it must not run in the generated path either. | ||
| val input = BoundReference(0, IntegerType, nullable = true) |
There was a problem hiding this comment.
Nit: this test never binds, evaluates or generates its child, since doGenCode throws before generating any child code, so it does not need a second BoundReference. Using a, as the SPARK-59121 test does for the same reduced shape (L101-104), removes the duplicate and checks "generates no code" more strictly: a is unevaluable, so a change that generated the child's code first would fail here.
Could we use a here?
There was a problem hiding this comment.
Removed with the doGenCode change in 913f317.
| } | ||
|
|
||
| /** `bucket` over integers. The functions below differ only in how Spark calls them. */ | ||
| private trait IntBucketFunction extends ScalarFunction[Int] { |
There was a problem hiding this comment.
Nit: these fixtures do not override canonicalName(), whose default returns a new random UUID on every call (BoundFunction.java:104-110), so two transforms over the same fixture object get different functionIds, and isSameFunction/hasSameReducedKeys treat them as different functions. Nothing here depends on it today, but the fixtures are top-level, so other tests in this package can pick them up, and the reduced-keys test's message embeds a random UUID. NamedFunction above has a stable canonical name for the same reason.
Could we add override def canonicalName(): String = "test.bucket" here?
There was a problem hiding this comment.
Removed with the doGenCode change in 913f317.
| |i.id, i.name, i.arrive_time | ||
| |FROM testcat.ns.$items i JOIN testcat.ns.$purchases p ON $joinCondition | ||
| |""".stripMargin) | ||
| checkAnswer(df, Seq( |
There was a problem hiding this comment.
Minor: neither new test can tell whether the merge orders rows by the value the transform's generated code computes, because checkAnswer ignores order. In the reported-ordering test, the only rows the merge compares are (1, 'aa', 2021) and (1, 'ab', 2022), and name decides before years(arrive_time), so the transform's code is compiled but never run. In the key-derived test it runs, but the scan already sorts the splits by the full key, so a transform that returned a constant or a wrong value would still pass.
Could the reported-ordering test make id and name tie, insert the later year first in a separate INSERT, and assert the merged order, e.g. with merging.head.executeCollect()?
There was a problem hiding this comment.
After 913f317 the merge no longer compares rows by the transform, so there is no transform order left to check. The tests now check the ordering the merge runs over, [id, name] and [id].
| val merging = collectAllGroupPartitions(df.queryExecution.executedPlan) | ||
| .filter(_.enableSortedMerge) | ||
| assert(merging.length == 1, "expected one k-way merge") | ||
| assert(merging.head.child.outputOrdering.exists(_.child.isInstanceOf[TransformExpression]), |
There was a problem hiding this comment.
No requirement can contain a transform, so the merge evaluates this key only as a tie-breaker nobody consumes. EnsureRequirements enables the merge only when the merged ordering satisfies the parent's requirement, SortOrder.orderingSatisfies is prefix-based with semantic equality, and no requiredChildOrdering can hold a TransformExpression: SMJ, aggregate and window keys are query expressions, and the write path replaces transforms with their calls. So in these two tests the requirements are [id, name] and [id], and years(arrive_time) only breaks ties after them. In the key-derived case every comparison within a coalesced group ties on id, so each one calls the connector function twice to compute a value that is constant per split. Cutting kWayMergeOrdering (GroupPartitionsExec.scala:275-276) and the merge branch of outputOrdering (GroupPartitionsExec.scala:374-377) before the first sort key that holds a transform would leave every merge decision unchanged, avoid those calls, and also avoid the failure for a transform Spark cannot call (see the comment on TransformExpression.scala:194).
Could we consider that, either instead of generating the transform's code here or in addition to it?
There was a problem hiding this comment.
Thanks, this is the simpler fix. I did it in the scan in 913f317, so every consumer of its ordering sees it, not only GroupPartitionsExec:
- in a reported ordering, a transform ends the leading run, like a pruned column;
- a transform key is dropped, while a sort order on another key still holds after it, since each partition holds a single key;
- the derived ordering leaves the transform keys out.
| // The scan reports no ordering, so it derives [id, years(arrive_time)] from its keys. The join | ||
| // on id projects the keys to id. With preserveKeyOrderingOnCoalesce off, the coalesced | ||
| // partitions keep no ordering on id unless they are merged. | ||
| createTable(items, itemsColumns, Array(identity("id"), years("arrive_time"))) |
There was a problem hiding this comment.
Minor: both new tests use years, which binds to YearsFunction with a magic invoke, so no end-to-end test runs a produceResult-only transform (the ApplyFunctionExpression route that the new comment in TransformExpression.scala singles out) through the merge's generated comparator; the catalyst test reaches that route only through projections. The suite already registers UnboundBucketFunction, whose BucketFunction has only produceResult and matches items.id (LongType).
Could we add a variant with a produceResult-only key, e.g. Array(identity("id"), bucket(4, "id")), so that the comparator evaluates it? Adding bucket to the first test's reported ordering would only compile it, since name decides first.
There was a problem hiding this comment.
After 913f317 the merge never calls a transform's function, so the way Spark calls it does not matter there. The derived test uses days, which Spark cannot call at all.
TransformExpression generate the code of its bound function…put ordering Reverts the doGenCode delegation. No operator requires an ordering over a transform, so the k-way merge no longer compares rows by one. The k-way merge comparator, now LazyRowOrdering, builds its comparator with RowOrdering.create, so it falls back to interpreted evaluation when code generation fails.
9170b23 to
913f317
Compare
|
Thank you for the review, @dongjoon-hyun. I followed items 5 and 3 in 913f317. The scan now drops every sort order that holds a partition transform, and the |
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for addressing items 1-12, @peter-toth. I reviewed the latest commit, 913f317, and left 9 inline comments. Here is a summary, ordered by priority, continuing the numbering from items 1-12.
Backports
KeyGroupedPartitioningSuiteL6662: thet3case and the derived k-way merge test rely onpartitionKeyOrderingdefaulting to true, which holds only on master and branch-4.x, so they fail if this is backported to branch-4.3 or branch-4.2.
Correctness
DataSourceV2ScanExecBaseL151: the write path sorts by the transform's function call only when Spark can call the function; this also corrects my item 5.
API and design
DataSourceV2ScanExecBaseL161: a sort order on a transform that is itself a partition key ends the leading run although it is constant within each partition, so the sort orders after it that still hold are dropped.
Docs and comments
DataSourceV2ScanExecBaseL169: the 4.4 migration note, thepartitionKeyOrderingconfig doc and the tuning guide still say the scan reports itself sorted by its partition key expressions.DataSourceV2ScanExecBaseL168: the PR description does not say that the failure exists since 4.2.0, nor that the new ordering can turn a hash aggregate into a sort aggregate.DataSourceV2ScanExecBaseL164: theBoundFunctionJavadoc still lists "retaining a reported ordering that matches the partitioning" as a benefit ofequals, which no longer applies.DataSourceV2ScanExecBaseL148: the coalescing-branch comment at GroupPartitionsExec.scala:391-395 still says the child's ordering is in sync with the partitioning.DataSourceV2ScanExecBaseL152: two sentences of the new scaladoc repeat earlier ones, and "even constant" suggests a property only transforms have.
Tests
KeyGroupedPartitioningSuiteL6591: the new tests keep the transform's column in the scan output only through their SELECT lists; with it pruned, they would pass on master.
Item 13 matters most if this goes to branch-4.3 or branch-4.2. Item 14 needs only a wording change, and item 15 could be a follow-up JIRA.
| "(1, 10.0, cast('2021-01-01' as timestamp)), " + | ||
| "(2, 20.0, cast('2021-01-01' as timestamp))") | ||
|
|
||
| withSQLConf( |
There was a problem hiding this comment.
These tests pass only where spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled defaults to true, which is master and branch-4.x since SPARK-59396, so they fail if this fix is backported to branch-4.3 or branch-4.2. Both branches default the conf to false. SPARK-58324, SPARK-59279, SPARK-59899 and SPARK-59905 all went to both branches, and this suite applies cleanly to them. With the conf off, t3 (L6588) reports no ordering, so outputOrdering is empty and L6593 fails. Here the scan derives no ordering, so kWayMergeIsFeasible is false, a SortExec replaces the merge, and assert(merging.length == 1, ...) at L6624 fails. The SPARK-56241 tests in this suite set the conf explicitly (L6342, L6388). On branch-4.2, this test also needs V2_BUCKETING_ALLOW_JOIN_KEYS_SUBSET_OF_PARTITION_KEYS, but that one fails at compile time.
Could we wrap the loop at L6589-6596 in withSQLConf(SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "true") and add the same entry here?
There was a problem hiding this comment.
Done in 6eca5ef. The shape test and the derived merge test now set partitionKeyOrdering explicitly. The 4.2 backport will use the conf's 4.2 name for the subset-keys setting.
| * Either way, a sort order that holds a partition transform is dropped, even on a partition key. | ||
| * In a reported ordering it also ends the leading run, like a sort order over a pruned column. | ||
| * This loses nothing today, since no operator requires an ordering over a transform. The write | ||
| * path sorts by the transform's function call instead. Dropping it is also a safe way to handle |
There was a problem hiding this comment.
Minor: a correction to my item 5, which said that "the write path replaces transforms with their calls". That holds only when Spark can call the function, and so does "The write path sorts by the transform's function call instead" here. resolvedFunction is defined only for a ScalarFunction (TransformExpression.scala:177-182), and resolveTransformExpression (DistributionAndOrderingUtils.scala:97-99) keeps any other transform in the write's local sort. Take a function f that binds to a non-scalar BoundFunction with a semantic equals, a table t1 that reports the ordering [f(id)], and a table t2 that requires [f(id)] with an unspecified distribution. Before this PR, INSERT INTO t2 SELECT * FROM t1 worked, because RemoveRedundantSorts (RemoveRedundantSorts.scala:52-55) dropped the write's SortExec against the scan's ordering. Now the scan drops that sort order, so the SortExec stays and fails to evaluate f. No connector I know binds a transform to a non-scalar function, and such a write already failed from any other source, so I think only the wording needs to change.
Could L150-151 say "The write path sorts by the transform's function call instead, when Spark can call the function"?
There was a problem hiding this comment.
Thanks, reworded in 6eca5ef: "The write path sorts by the transform's function call instead, when Spark can call the function."
| case (Some(o), p) => | ||
| val (prefix, rest) = o.span(_.references.subsetOf(outputSet)) | ||
| val (prefix, rest) = | ||
| o.span(order => order.references.subsetOf(outputSet) && !holdsTransform(order.child)) |
There was a problem hiding this comment.
Minor: a sort order on a transform that is itself a partition key ends the leading run here, although it is constant within each partition, which is the reason L142-143 give for keeping the sort orders on a partition key. So the sort orders after it that still hold are dropped. With keys [days(ts)] and a reported [days(ts), id], the scan reports [] instead of [id]. With keys [years(ts), id] and a reported [years(ts), id, name], it reports [id] instead of [id, name]. A reported [days(ts)] over keys [days(ts), id] also gives [], while no report would derive [id]; V2ScanPartitioningAndOrdering.scala:98-101 uses None for a dropped report for that reason. This is not a regression, since none of these orderings satisfied anything before, and it comes from the cut "before the first sort key that holds a transform" that I suggested in item 5. Matching a transform sort order to a transform key would need isSameFunction plus semantically equal children rather than ExpressionSet, since the two are bound separately and BoundFunction.equals is only a SHOULD (BoundFunction.java:112-124).
Could we skip a sort order on a partition key instead of ending the run there, e.g. prefix ++ rest.filter(o => isKey(o) && usable(o)) ++ rest.takeWhile(o => isKey(o) || usable(o)).filterNot(isKey) with usable being the span predicate, and fall back to the derived ordering when nothing of a non-empty report survives? A follow-up JIRA would be fine too.
| } | ||
| case (_, k: KeyedPartitioning) if conf.v2BucketingPartitionKeyOrderingEnabled => | ||
| k.expressions.map(SortOrder(_, Ascending)) | ||
| k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending)) |
There was a problem hiding this comment.
Minor: with transform keys left out here, some docs now overstate what the derived ordering gives. The 4.4 migration note (docs/sql-migration-guide.md:45) says the scan "now also reports itself sorted by its partition key expressions", and the partitionKeyOrdering doc (SQLConf.scala:2591-2595) and docs/sql-performance-tuning.md:711 say the same. A table partitioned only by transforms, such as days(ts) or bucket(8, id), now derives no ordering, so the conf does nothing for it, and keys [years(ts), id] derive [id]. SupportsReportOrdering.java:36-42 could also tell connector authors that a sort order over a partition transform is not used.
Could we say in those docs that partition transforms are left out, e.g. "sorted by those of its partition key expressions that are columns (a partition transform such as days(ts) or bucket(8, id) is left out)"?
There was a problem hiding this comment.
Done in 6eca5ef, in the conf doc, the tuning guide, the 4.4 migration note and the SupportsReportOrdering Javadoc.
| prefix ++ rest.filter(order => keyExprs.contains(order.child)) | ||
| case _ => prefix | ||
| } | ||
| case (_, k: KeyedPartitioning) if conf.v2BucketingPartitionKeyOrderingEnabled => |
There was a problem hiding this comment.
Minor: the PR description does not say since when the failure exists or where the ordering change is on by default. The k-way merge (SPARK-55715) and the key-derived ordering (SPARK-56241) are both in 4.2.0, so the failure exists since 4.2.0 with spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled on, while this conf defaults to true only on master and branch-4.x. The new [id] can also turn a hash aggregate into a sort aggregate. spark.sql.execution.replaceHashWithSortAgg defaults to true on master, branch-4.x and branch-4.3, and ReplaceHashWithSortAgg.scala:58-65 switches when the child's ordering satisfies the grouping, so SELECT id, max(ts) FROM t GROUP BY id over keys [years(ts), id] now plans a SortAggregateExec. The template asks whether a change is user-facing compared to the released versions.
Could we add both to the user-facing section?
There was a problem hiding this comment.
Added both to the user-facing section. I measured the sort aggregate: SELECT id, max(ts) FROM t GROUP BY id over the keys [years(ts), id] plans its partial aggregate as a SortAggregateExec with this PR, and as a HashAggregateExec without it. One correction: preserveOrderingOnCoalesce is off by default on every branch. So the description says that the failure needs it, and that the derived case also needs partitionKeyOrdering, which defaults to true only on master and branch-4.x.
| p match { | ||
| case k: KeyedPartitioning if rest.nonEmpty => | ||
| val keyExprs = ExpressionSet(k.expressions) | ||
| val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform)) |
There was a problem hiding this comment.
Nit: BoundFunction.java:121 still lists "retaining a reported ordering that matches the partitioning" among the benefits of a semantic equals. That refers to this keyExprs.contains match and the one at GroupPartitionsExec.scala:397-398, but keyExprs now holds no transform, so whether a reported ordering is retained no longer depends on equals.
Could we drop that bullet, unless item 15 brings back an equals-based match?
There was a problem hiding this comment.
Dropped in 6eca5ef, but I put a different bullet in its place. PlanMerger compares the logical ordering, which keeps the transform, so merging two scans that report an ordering over the same transform still depends on a semantic equals.
| * is on, each partition contains rows where the key expressions evaluate to a single constant | ||
| * value, so the data is trivially sorted by those expressions within the partition. | ||
| * | ||
| * Either way, a sort order that holds a partition transform is dropped, even on a partition key. |
There was a problem hiding this comment.
Nit: this makes the coalescing-branch comment at GroupPartitionsExec.scala:391-395 stale. It says the child's outputOrdering "should already be in sync with the partitioning (either reported by the source or derived from it in DataSourceV2ScanExecBase)", but this scan now leaves out every transform key, so the transform entries of keyExprs there never match. For example, keys [years(ts)] with preserveKeyOrderingOnCoalesce on reported [years(ts)] before and report [] now. That file is already in this PR.
Could that comment say the child's ordering is in sync with the partitioning's column keys, since this scan drops every sort order over a partition transform?
| * In a reported ordering it also ends the leading run, like a sort order over a pruned column. | ||
| * This loses nothing today, since no operator requires an ordering over a transform. The write | ||
| * path sorts by the transform's function call instead. Dropping it is also a safe way to handle | ||
| * a transform Spark cannot evaluate. Keeping it would only add comparisons nobody uses. A sort |
There was a problem hiding this comment.
Nit: "Keeping it would only add comparisons nobody uses." repeats the argument of L150 two sentences later, and "A sort order that Spark derives from a transform key is even constant within each partition." repeats L144-146: every derived sort order is constant within a partition, identity keys included, so "even" reads as a reason specific to transforms, which it is not. "only" also understates it: a kept transform ahead of a key hides that key from prefix matching, as [years(ts), id] vs [id] shows.
Could L150-154 be shortened to something like "No operator requires an ordering over a transform (the write path sorts by the transform's function call when Spark can call it), so keeping it would add comparisons nobody uses, and dropping it also handles a transform Spark cannot evaluate. Revisit this if such a requirement appears."?
There was a problem hiding this comment.
Shortened in 6eca5ef. "only" and "even constant" are gone, and the write-path sentence now has the qualifier from item 14.
| createTable(table3, columns, Array(days("ts"), identity("id"))) | ||
| Seq(table1, table2, table3).foreach { table => | ||
| sql(s"INSERT INTO testcat.ns.$table VALUES (1, 'aa', cast('2020-01-01' as timestamp))") | ||
| val df = sql(s"SELECT id, data, ts FROM testcat.ns.$table") |
There was a problem hiding this comment.
Nit: these tests keep the transform's column in the scan output only through their SELECT lists, here and i.arrive_time at L6615. With it pruned, the pruned-column handling from SPARK-59899 already gives the same orderings on master: [id] for all three tables here, and [id, name] and [id] in the merge tests. So a later cleanup of a SELECT list would silently stop these tests from guarding the fix, and this one has no checkAnswer to notice. The SPARK-59899 test asserts the opposite precondition (L3910-3911).
Could we assert it here too, e.g. assert(scan.output.exists(_.name == "ts"), s"test setup: ts stays in the output of $table"), and mention it in the scaladoc of checkKWayMergeBeforeTransform?
There was a problem hiding this comment.
Added the setup asserts in 6eca5ef, and said in the scaladoc of checkKWayMergeBeforeTransform why the column has to stay.
There was a problem hiding this comment.
Thank you, @peter-toth. Following review 5445587012, I left 3 more inline comments on the latest commit, 913f317, continuing the numbering from items 13-21.
Correctness
GroupPartitionsExecL337: the merge re-readschild.outputOrderingat execution, and its derived part followspartitionKeyOrdering.enabledlive, so turning that conf off after planning leaves a planned merge with an empty ordering; this predates the PR and fits a follow-up like SPARK-59279.
Tests
GroupPartitionsExecSuiteL171: the fallback test passes whether or not code generation failed.
Cleanup
KeyGroupedPartitioningSuiteL6583:sort(expr, SortDirection.ASCENDING)already defaults toNULLS_FIRST.
Process
- SPARK-59995: its summary probably still describes the reverted codegen change; could you retitle it to match the new PR title? Otherwise
dev/merge_spark_pr.pywarns about a low title similarity at merge time.
Item 22 matters most, and it can be a follow-up JIRA.
| } else if (usesSortedMerge) { | ||
| val partitionCoalescer = new GroupedPartitionCoalescer(groupedPartitions.map(_._2)) | ||
| val rowOrdering = new LazyCodeGenOrdering(kWayMergeOrdering, child.output) | ||
| val rowOrdering = new LazyRowOrdering(kWayMergeOrdering, child.output) |
There was a problem hiding this comment.
kWayMergeOrdering is read from the child again at execution time, and the scan's derived part of it follows partitionKeyOrdering.enabled live, so a merge planned over a derived ordering can run with an empty one. DataSourceV2ScanExecBase.outputOrdering derives the key ordering only while spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled is on (DataSourceV2ScanExecBase.scala:168), and both this line and kWayMergeIsFeasible (L247-248) read child.outputOrdering when the node executes. For example, take a table identity-partitioned by (a, b) that reports no ordering, with allowKeysSubsetOfPartitionKeys and preserveOrderingOnCoalesce on, and plan SELECT a, b, c, row_number() OVER (PARTITION BY a ORDER BY b) FROM t: the plan has GroupPartitionsExec(SortedMerge: true) under the window and no SortExec. If the conf is turned off before collect(), the node can still merge, since its usesSortedMerge may already have been evaluated during planning, but kWayMergeOrdering is now empty, so LazyRowOrdering(Nil) treats every row as equal and the window gets its rows out of order. This predates the PR (SPARK-56241) and needs two non-default confs plus a conf change after planning, so it is the same kind of issue as SPARK-59279 rather than a blocker here.
Could we file a follow-up JIRA to freeze the merge ordering when tryEnableSortedMerge decides on the merge, as SPARK-59279 did for the merge config?
| Seq(SortOrder(TransformExpression(YearsFunction, Seq(ts)), Ascending)), Seq(ts)) | ||
| val serializer = new JavaSerializer(new SparkConf()).newInstance() | ||
| val shipped = serializer.deserialize[LazyRowOrdering](serializer.serialize(ordering)) | ||
| withSQLConf(SQLConf.CODEGEN_FACTORY_MODE.key -> CodegenObjectFactoryMode.FALLBACK.toString) { |
There was a problem hiding this comment.
Nit: these assertions hold whether or not code generation failed, so this pins the fallback only while TransformExpression.doGenCode throws. If a transform generates code again, as in the round-1 version of this PR, the test would still pass without any fallback. And since the scan no longer lets a transform reach the merge, the transform here only stands in for a sort key without generated code, which the comment at L163-165 could say.
Could we also check that a fresh deserialized copy fails under the session's CODEGEN_ONLY, e.g. intercept[SparkException] with "Cannot generate code for expression", so the test fails if the premise changes?
There was a problem hiding this comment.
Added in 6eca5ef. A fresh copy now fails under the session's CODEGEN_ONLY with "Cannot generate code for expression", and the comment says that the transform only stands in for a sort key without generated code.
| val table2 = "transform_order_t2" | ||
| val table3 = "transform_order_t3" | ||
| def asc(expr: Expression): SortOrder = | ||
| sort(expr, SortDirection.ASCENDING, NullOrdering.NULLS_FIRST) |
There was a problem hiding this comment.
Nit: the two-argument sort(expr, SortDirection.ASCENDING) already uses NullOrdering.NULLS_FIRST, the default null ordering of SortDirection.ASCENDING, and this suite already uses that form (e.g. L4043).
Could this be def asc(expr: Expression): SortOrder = sort(expr, SortDirection.ASCENDING)?
There was a problem hiding this comment.
Done in 6eca5ef, and in the reported merge test's sort orders too.
|
This morning, GitHub had some outage and it seems that there exists some glitches still. I revised the above review comments due to that, @peter-toth . |
Pin partitionKeyOrdering in the tests for the maintenance branches, keep the transform's column in the scan output on purpose, check the codegen premise of the fallback test, and say in the docs that the derived ordering leaves partition transforms out.
|
Thank you for the second round, @dongjoon-hyun. I addressed items 13, 14, 16-21, 23 and 24 in 6eca5ef, and filed SPARK-60043 for item 15 and SPARK-60044 for item 22. I also retitled SPARK-59995 to match the PR title (item 25). |
dongjoon-hyun
left a comment
There was a problem hiding this comment.
+1, LGTM (Pending CIs).
| * including columns pruned from the scan output. If any of them cannot be resolved, Spark | ||
| * ignores the whole reported ordering and logs a warning. | ||
| * <p> | ||
| * Spark does not use a sort order over a transform such as {@code days(ts)}, whether or not it is |
There was a problem hiding this comment.
Minor: this new public Javadoc says more than holds.
- The sort orders on partition keys after a transform are kept only when Spark keeps the reported
KeyGroupedPartitioning, i.e. withspark.sql.sources.v2.bucketing.enabledon and every split implementingHasPartitionKey. Otherwise they are dropped too (DataSourceV2ScanExecBase.scala:166). - "Spark does not use a sort order over a transform" holds for physical planning only.
PlanMerger.combineRequiredOrderingandmergeDegradesReporting(PlanMerger.scala:967-1001) still compare the reported ordering with the transform in it, which is what the newBoundFunctionbullet says. - SPARK-60043 plans to change "every one after it", so this public interface doc would have to change again.
Could this be less committal, e.g. "Spark currently does not exploit a sort order over a transform such as {@code days(ts)} when planning, nor the sort orders after it, except those on partition keys that are not transforms when Spark uses the reported {@code KeyGroupedPartitioning}"?
There was a problem hiding this comment.
Thanks, reworded in f77f37e: "Spark currently does not rely on a sort order over a transform such as days(ts) to avoid a sort, whether or not it is a partition key. Nor does it rely on the sort orders after it. The exception is a sort order on a partition key that is not a transform, when Spark uses the reported KeyGroupedPartitioning." I used "rely on ... to avoid a sort" rather than "when planning", since PlanMerger still compares these sort orders (SPARK-60083).
| * <li>accepting a valid query that selects and groups by the same scalar function call</li> | ||
| * <li>keeping a union's keyed partitioning</li> | ||
| * <li>retaining a reported ordering that matches the partitioning</li> | ||
| * <li>merging scans that report an ordering over the same transform</li> |
There was a problem hiding this comment.
Minor: the scan merge depends on equals for a reported partitioning over a transform too, not only for an ordering. PlanMerger.combineRequiredKeyGroupedPartitioning (PlanMerger.scala:955-961) compares the two inputs' reported keys by canonical equality, and the SPARK-58769 test in MergeSubplansSuite ("decline the merge when the reported transform is not comparable") pins exactly that for bucket(4, a). That is the more common case, and after this PR it is the one that still matters for the physical plan.
Could this bullet say "merging scans that report a partitioning or an ordering over the same transform"?
| * becomes a real requirement. | ||
| */ | ||
| override def outputOrdering: Seq[SortOrder] = { | ||
| def holdsTransform(e: Expression): Boolean = e.exists(_.isInstanceOf[TransformExpression]) |
There was a problem hiding this comment.
Minor: the transform sort orders are dropped only here, in the physical outputOrdering. The logical DataSourceV2ScanRelation.ordering keeps them, so PlanMerger (combineRequiredOrdering and mergeDegradesReporting, PlanMerger.scala:967-1001) and DataSourceV2ScanRelation.doCanonicalize (DataSourceV2Relation.scala:414-418) still gate scan merging and subplan deduplication on sort orders that no plan uses any more. For example, take a SCAN_MERGING table that reports [years(ts)] and whose years function does not override equals. Two scalar subqueries over it bind years separately, so neither ordering satisfies the other and the merge is declined (unless spark.sql.optimizer.mergeSubplans.dsv2ScanMerge.orderingDegradation.enabled is on), although both physical scans now report []. The new BoundFunction bullet documents this cost rather than removing it.
Could we file a follow-up JIRA to leave the transform sort orders out of those logical comparisons as well?
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partition.filter.enabled` defaults to `true`. In a storage-partitioned join, the partition key groups pushed down to both sides may now be narrowed to those that can produce output for the join type, instead of always taking the union of the two sides' groups: an inner or semi join may keep only the groups present on both sides, a join that keeps or tests every left row (left outer, left anti, left single, existence) keeps the left side's groups, a right outer join keeps the right side's, and a full outer join is unaffected; a CROSS join carrying an equality condition is treated like an inner join. The narrowing is skipped when either side's partitioning may contain unknown partition keys, as after a shuffle on one side, so such a join keeps the full union. A group that is dropped is not scanned at all, so a query may read fewer files. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partition.filter.enabled` to `false`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate whose grouping is exactly the partition key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate that groups by those key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. Partition transforms such as `days(ts)` or `bucket(8, id)` are left out of that ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` defaults to `true`. When `GroupPartitionsExec` merges several input partitions that share one partition key value into a single output partition, it now keeps sort orders over the partition key expressions in the ordering it reports, since those expressions are constant within the merged partition. Orders over other columns are still dropped, as the concatenation invalidates them, and a join that reduced the partition keys onto a common key space (see `spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled`) reports no order, since the merged partitions then share only the reduced key. This can remove a `Sort` downstream of the merge. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled` to `false`. |
There was a problem hiding this comment.
Minor: like the partitionKeyOrdering docs in item 16, the preserveKeyOrderingOnCoalesce docs now overstate what is kept: this note ("keeps sort orders over the partition key expressions"), the conf doc (SQLConf.scala:2605-2610) and docs/sql-performance-tuning.md:719. Sort orders over transform keys no longer reach GroupPartitionsExec. With keys [days(ts)] and a reported [days(ts)], a coalescing GroupPartitionsExec reported [days(ts)] before this PR and reports nothing now.
Could we add the same "partition transforms such as days(ts) or bucket(8, id) are left out" to those three places?
There was a problem hiding this comment.
Done in f77f37e, in all three places. The conf doc also called the reported ordering "key-derived", which I fixed while there.
| } | ||
| } | ||
|
|
||
| test("SPARK-59995: a scan's output ordering drops sort orders that hold a partition transform") { |
There was a problem hiding this comment.
Nit: the sort-aggregate change that the PR description lists as user-facing has no test: SELECT id, max(ts) FROM t GROUP BY id over the keys [years(ts), id] now plans a partial SortAggregateExec. A later change to ReplaceHashWithSortAgg or to this ordering could flip it silently.
Could this test also check, for t2, that GROUP BY id plans a SortAggregateExec, with spark.sql.execution.replaceHashWithSortAgg set explicitly for the backports?
There was a problem hiding this comment.
Added in f77f37e. The shape test now checks that GROUP BY id plans a SortAggregateExec for each table, with replaceHashWithSortAgg set explicitly.
| - Since Spark 4.4, `/*+ ... */` inside a bracketed comment is parsed as a nested comment, so its closing `*/` no longer closes the outer comment. For example, `/* note /*+ x */ SELECT 1` previously returned `1`, but now raises `UNCLOSED_BRACKETED_COMMENT`. Close the outer comment explicitly, for example `/* note /*+ x */ */ SELECT 1`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partition.filter.enabled` defaults to `true`. In a storage-partitioned join, the partition key groups pushed down to both sides may now be narrowed to those that can produce output for the join type, instead of always taking the union of the two sides' groups: an inner or semi join may keep only the groups present on both sides, a join that keeps or tests every left row (left outer, left anti, left single, existence) keeps the left side's groups, a right outer join keeps the right side's, and a full outer join is unaffected; a CROSS join carrying an equality condition is treated like an inner join. The narrowing is skipped when either side's partitioning may contain unknown partition keys, as after a shuffle on one side, so such a join keeps the full union. A group that is dropped is not scanned at all, so a query may read fewer files. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partition.filter.enabled` to `false`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate whose grouping is exactly the partition key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate that groups by those key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. Partition transforms such as `days(ts)` or `bucket(8, id)` are left out of that ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. |
There was a problem hiding this comment.
Nit: no migration note covers the reported-ordering half of this change, which applies under any conf. A source that reports [years(ts), id] over the keys [years(ts), id] now gets [id], which can drop a Sort and turn a hash aggregate into a sort aggregate, as the PR description says. This note covers only the derived ordering, and its "To restore the previous behavior" sentence does not apply to the reported one. It only removes sorts, so this is low priority.
Could we add a sentence for it, or would you prefer to leave it to the PR description?
There was a problem hiding this comment.
I'd leave it to the PR description. The change only removes work, and this fix goes back to branch-4.3 and branch-4.2, where a 4.4 migration note would not apply.
| - Since Spark 4.4, `/*+ ... */` inside a bracketed comment is parsed as a nested comment, so its closing `*/` no longer closes the outer comment. For example, `/* note /*+ x */ SELECT 1` previously returned `1`, but now raises `UNCLOSED_BRACKETED_COMMENT`. Close the outer comment explicitly, for example `/* note /*+ x */ */ SELECT 1`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partition.filter.enabled` defaults to `true`. In a storage-partitioned join, the partition key groups pushed down to both sides may now be narrowed to those that can produce output for the join type, instead of always taking the union of the two sides' groups: an inner or semi join may keep only the groups present on both sides, a join that keeps or tests every left row (left outer, left anti, left single, existence) keeps the left side's groups, a right outer join keeps the right side's, and a full outer join is unaffected; a CROSS join carrying an equality condition is treated like an inner join. The narrowing is skipped when either side's partitioning may contain unknown partition keys, as after a shuffle on one side, so such a join keeps the full union. A group that is dropped is not scanned at all, so a query may read fewer files. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partition.filter.enabled` to `false`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate whose grouping is exactly the partition key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. | ||
| - Since Spark 4.4, `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` defaults to `true`. A V2 scan that reports a keyed partitioning but no explicit ordering now also reports itself sorted by its partition key expressions, because every row of such a partition evaluates them to the same value. Spark can then drop a `Sort` it would otherwise place above the scan. An aggregate that groups by those key expressions may also be planned as a sort aggregate rather than a hash aggregate, because `spark.sql.execution.replaceHashWithSortAgg` keys off the reported ordering. Partition transforms such as `days(ts)` or `bucket(8, id)` are left out of that ordering. To restore the previous behavior, set `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` to `false`. |
There was a problem hiding this comment.
Nit: "An aggregate that groups by those key expressions" now comes before the sentence that leaves the transforms out, so it reads as if grouping by a transform key could switch. Also, "exactly" is gone, although ReplaceHashWithSortAgg matches the ordering by prefix. With keys [a, id], GROUP BY id or GROUP BY id, a still plans a hash aggregate, and with keys [days(ts)] no aggregate can switch.
Could the transform sentence come first, and the aggregate one say "groups by a leading run of those key columns, in key order"?
There was a problem hiding this comment.
Done in f77f37e. The transform sentence comes first, and the aggregate one now says "groups by a leading run of the other key expressions, in key order".
| } | ||
|
|
||
| test("SPARK-59995: k-way merge over an ordering derived from a partition transform key") { | ||
| // The scan reports no ordering, so it derives [id, days(arrive_time)] from its keys and keeps |
There was a problem hiding this comment.
Nit: with this PR the scan never derives [id, days(arrive_time)]. k.expressions.filterNot(holdsTransform) (DataSourceV2ScanExecBase.scala:169) leaves the transform key out of the derivation itself, so nothing is cut down to [id].
Could this say "so it derives [id] from its keys, leaving out days(arrive_time)"?
Narrow the SupportsReportOrdering and BoundFunction Javadocs, say in the preserveKeyOrderingOnCoalesce docs that transform keys are left out, and check the sort aggregate in the scan ordering test.
|
Thank you for the approval, @dongjoon-hyun. I addressed the new comments in f77f37e, and filed SPARK-60083 for the logical comparisons. |
…rm-expression-codegen # Conflicts: # sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/v2/GroupPartitionsExecSuite.scala
…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 #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) Signed-off-by: Peter Toth <peter.toth@gmail.com>
|
Thank you @dongjoon-hyun for the review. |
…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>
…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>
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)