Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions docs/sql-migration-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,8 @@ license: |
- Since Spark 4.4, when `array_repeat` or `array_insert` is asked to build an array larger than the maximum supported array length, generated code raises the same error as interpreted evaluation. `array_repeat` now fails with `COLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER` instead of the internal error `_LEGACY_ERROR_TEMP_2176`, and `array_insert` fails with `COLLECTION_SIZE_LIMIT_EXCEEDED.FUNCTION` instead of `COLLECTION_SIZE_LIMIT_EXCEEDED.PARAMETER`, which named a `count` parameter that `array_insert` does not have. Both functions raise an error under exactly the same conditions as before; only the reported error condition changes.
- 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.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`.
- 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. Partition transforms such as `days(ts)` or `bucket(8, id)` are left out of that ordering. This lets Spark drop a `Sort` it would otherwise place above the scan. An aggregate that groups by a leading run of the other key expressions, in key order, may also be planned as a sort aggregate rather than a hash aggregate. This is 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.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. Sort orders over partition transforms such as `days(ts)` or `bucket(8, id)` are left out. 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`.

## Upgrading from Spark SQL 4.2 to 4.3

Expand Down
4 changes: 2 additions & 2 deletions docs/sql-performance-tuning.md
Original file line number Diff line number Diff line change
Expand Up @@ -708,15 +708,15 @@ The following SQL properties enable Storage Partition Join in different join que
<td><code>spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled</code></td>
<td>true</td>
<td>
When enabled, Spark derives the output ordering of a V2 scan from its partition key expressions, if the source reports a keyed partitioning but no explicit ordering, or an ordering that Spark ignores because it references a column that cannot be resolved. All rows of such a partition share one key value, so the partition is trivially sorted by those expressions, and a sort Spark would otherwise add becomes unnecessary. This config requires <code>spark.sql.sources.v2.bucketing.enabled</code> to be true.
When enabled, Spark derives the output ordering of a V2 scan from its partition key expressions, if the source reports a keyed partitioning but no explicit ordering, or an ordering that Spark ignores because it references a column that cannot be resolved. All rows of such a partition share one key value, so the partition is trivially sorted by those expressions, and a sort Spark would otherwise add becomes unnecessary. Partition transforms such as <code>days(ts)</code> or <code>bucket(8, id)</code> are left out of the ordering. This config requires <code>spark.sql.sources.v2.bucketing.enabled</code> to be true.
</td>
<td>4.2.0</td>
</tr>
<tr>
<td><code>spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled</code></td>
<td>true</td>
<td>
When enabled, <code>GroupPartitionsExec</code> reports sort orders over partition key expressions after coalescing several input partitions into one. The merged partitions share the same partition key value, so these orders still hold, while orders over other columns are lost by the concatenation. No order is reported when the join reduced the partition keys onto a common key space (see <code>spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled</code>), because the merged partitions then share only the reduced key. This config requires <code>spark.sql.sources.v2.bucketing.enabled</code> to be true.
When enabled, <code>GroupPartitionsExec</code> reports sort orders over partition key expressions after coalescing several input partitions into one. The merged partitions share the same partition key value, so these orders still hold, while orders over other columns are lost by the concatenation. Sort orders over partition transforms such as <code>days(ts)</code> or <code>bucket(8, id)</code> are left out. No order is reported when the join reduced the partition keys onto a common key space (see <code>spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled</code>), because the merged partitions then share only the reduced key. This config requires <code>spark.sql.sources.v2.bucketing.enabled</code> to be true.
</td>
<td>4.2.0</td>
</tr>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -118,7 +118,7 @@ default String canonicalName() {
* <ul>
* <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 a partitioning or an ordering over the same transform</li>
* <li>reusing identical bucketed scans</li>
* <li>recognizing two identical subplans</li>
* </ul>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,11 @@ public interface SupportsReportOrdering extends Scan {
* Spark resolves the column references in these sort orders against the table columns,
* 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 currently does not rely on a sort order over a transform such as {@code 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 {@code KeyGroupedPartitioning}.
*/
SortOrder[] outputOrdering();
}
Original file line number Diff line number Diff line change
Expand Up @@ -2655,7 +2655,8 @@ object SQLConf {
"a V2 data source that reports a KeyedPartitioning but does not report explicit ordering " +
"via SupportsReportOrdering, or reports one that Spark ignores because it references a " +
"column that cannot be resolved. Within a single partition all rows share the same key " +
s"value, so the data is trivially sorted by those expressions. Requires " +
"value, so the data is trivially sorted by those expressions. Partition transforms such " +
"as `days(ts)` or `bucket(8, id)` are left out of the ordering. Requires " +
s"${V2_BUCKETING_ENABLED.key} to be enabled.")
.version("4.2.0")
.withBindingPolicy(ConfigBindingPolicy.SESSION)
Expand All @@ -2667,9 +2668,10 @@ object SQLConf {
.doc("When enabled, Spark preserves sort orders over partition key expressions when " +
"GroupPartitionsExec coalesces multiple input partitions into one output partition. " +
"Because all merged partitions share the same partition key value, sort orders over " +
"those key expressions remain valid after the merge. This applies to both key-derived " +
"ordering (from SupportsReportOrdering) and ordering derived from " +
s"${V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key}. Requires " +
"those key expressions remain valid after the merge. This applies to both the ordering " +
"reported via SupportsReportOrdering and the ordering derived from " +
s"${V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key}. Sort orders over partition " +
"transforms such as `days(ts)` or `bucket(8, id)` are left out. Requires " +
s"${V2_BUCKETING_ENABLED.key} to be enabled.")
.version("4.2.0")
.withBindingPolicy(ConfigBindingPolicy.SESSION)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ package org.apache.spark.sql.execution.datasources.v2

import org.apache.spark.rdd.RDD
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, ExpressionSet, SortOrder}
import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, ExpressionSet, SortOrder, TransformExpression}
import org.apache.spark.sql.catalyst.plans.physical
import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning
import org.apache.spark.sql.catalyst.util.truncatedString
Expand Down Expand Up @@ -144,19 +144,29 @@ trait DataSourceV2ScanExecBase
* is a `KeyedPartitioning` and `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled`
* 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.

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.

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Reworded in 6eca5ef.

* In a reported ordering it also ends the leading run, like a sort order over a pruned column.
* Dropping it 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. Keeping such a sort order would add comparisons nobody uses. Dropping it also
* handles a transform Spark cannot evaluate. Revisit this if an ordering over a transform
* becomes a real requirement.
*/
override def outputOrdering: Seq[SortOrder] = {
def holdsTransform(e: Expression): Boolean = e.exists(_.isInstanceOf[TransformExpression])

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.

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed. I filed SPARK-60083 for it.

(ordering, outputPartitioning) match {
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))

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.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed. I filed SPARK-60043 for it.

p match {
case k: KeyedPartitioning if rest.nonEmpty =>
val keyExprs = ExpressionSet(k.expressions)
val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform))

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.

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

prefix ++ rest.filter(order => keyExprs.contains(order.child))
case _ => prefix
}
case (_, k: KeyedPartitioning) if conf.v2BucketingPartitionKeyOrderingEnabled =>

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.

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

k.expressions.map(SortOrder(_, Ascending))
k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending))

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.

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)"?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done in 6eca5ef, in the conf doc, the tuning guide, the 4.4 migration note and the SupportsReportOrdering Javadoc.

case _ => Seq.empty
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@ import org.apache.spark.{Partition, SparkException}
import org.apache.spark.rdd.{CoalescedRDD, PartitionCoalescer, PartitionGroup, RDD, SortedMergeCoalescedRDD}
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions._
import org.apache.spark.sql.catalyst.expressions.codegen.GenerateOrdering
import org.apache.spark.sql.catalyst.plans.QueryPlan
import org.apache.spark.sql.catalyst.plans.physical.{IdentityReducer, KeyedPartitioning, KeyLayout, KeyReducer, Partitioning, PartitioningCollection, REPLICATED_FOR_JOIN, UngroupingOrigin, UnknownPartitioning}
import org.apache.spark.sql.catalyst.util.{truncatedString, InternalRowComparableWrapper}
Expand Down Expand Up @@ -277,8 +276,8 @@ case class GroupPartitionsExec(
}

/**
* The ordering used by the k-way merge in [[SortedMergeCoalescedRDD]]. The generated comparator
* ([[GenerateOrdering]]) only needs each [[SortOrder]]'s sort key (child, direction, null
* The ordering used by the k-way merge in [[SortedMergeCoalescedRDD]]. The comparator that
* [[RowOrdering]] builds only needs each [[SortOrder]]'s sort key (child, direction, null
* ordering), so `sameOrderExpressions` -- planner-only metadata that would otherwise be
* serialized with the RDD in every task -- is dropped.
*/
Expand Down Expand Up @@ -345,7 +344,7 @@ case class GroupPartitionsExec(
sparkContext.emptyRDD
} else if (usesSortedMerge) {
val partitionCoalescer = new GroupedPartitionCoalescer(groupedPartitions.map(_._2))
val rowOrdering = new LazyCodeGenOrdering(kWayMergeOrdering, child.output)
val rowOrdering = new LazyRowOrdering(kWayMergeOrdering, child.output)

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.

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed. I filed SPARK-60044 for it.

new SortedMergeCoalescedRDD[InternalRow](
child.execute(),
groupedPartitions.size,
Expand Down Expand Up @@ -399,11 +398,9 @@ case class GroupPartitionsExec(
outputPartitioning match {
case p: Partitioning with Expression
if reducers.isEmpty && conf.v2BucketingPreserveKeyOrderingOnCoalesceEnabled =>
// Without reducers all merged partitions share the same original key value, so the key
// expressions remain constant within the output partition. The child's outputOrdering
// should already be in sync with the partitioning (either reported by the source or
// derived from it in DataSourceV2ScanExecBase), so we only need to keep the sort orders
// whose expression is a partition key expression -- all others are lost by concatenation.
// Without reducers all merged partitions share the same original key value, so the sort
// orders on key expressions still hold. The transform keys match nothing here, since
// `DataSourceV2ScanExecBase.outputOrdering` drops every sort order over a transform.
val keyedPartitionings = p.collect { case k: KeyedPartitioning => k }
val keyExprs = ExpressionSet(keyedPartitionings.flatMap(_.expressions))
child.outputOrdering.filter(order => keyExprs.contains(order.child))
Expand Down Expand Up @@ -820,15 +817,15 @@ class GroupedPartitionCoalescer(
}

/**
* A serializable [[Ordering]] for [[InternalRow]] that generates code-compiled comparison logic
* lazily on first use. The [[SortOrder]] expressions and output schema are serialized with the
* RDD; the generated comparator is rebuilt on the executor on first comparison via
* [[GenerateOrdering]].
* A serializable [[Ordering]] for [[InternalRow]]. The [[SortOrder]] expressions and output schema
* are serialized with the RDD. The comparator is built on the executor, at the first comparison.
* [[RowOrdering]] builds it, as for `SortExec`. It tries generated code first and falls back to
* interpreted evaluation when that fails.
*/
private class LazyCodeGenOrdering(
private class LazyRowOrdering(
sortOrders: Seq[SortOrder],
schema: Seq[Attribute]) extends Ordering[InternalRow] with Serializable {
@transient private lazy val generated: Ordering[InternalRow] =
GenerateOrdering.generate(sortOrders, schema)
override def compare(x: InternalRow, y: InternalRow): Int = generated.compare(x, y)
@transient private lazy val ordering: Ordering[InternalRow] =
RowOrdering.create(sortOrders, schema)
override def compare(x: InternalRow, y: InternalRow): Int = ordering.compare(x, y)
}
Loading