Skip to content

[SPARK-59981][SQL] Derive the partition-key ordering when a V2 scan's reported ordering keeps no sort order - #59257

Open
peter-toth wants to merge 2 commits into
apache:masterfrom
peter-toth:SPARK-59981-empty-reported-ordering
Open

peter-toth wants to merge 2 commits into
apache:masterfrom
peter-toth:SPARK-59981-empty-reported-ordering

Conversation

@peter-toth

@peter-toth peter-toth commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

DataSourceV2ScanExecBase.outputOrdering now derives an ordering from the partition key expressions whenever the reported ordering keeps no sort order. Before, it derived one only when the scan reported no ordering at all.

  • It still restricts the reported ordering to the scan output first, and drops the sort orders over a partition transform. If nothing is left, the output partitioning is a KeyedPartitioning and spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled is on, it returns the ordering over the partition keys.
  • Three cases lost the key ordering before:
    • an empty report, which V2ScanPartitioningAndOrdering stores as Some(Nil);
    • a report that starts with a sort order on a pruned column and has no sort order on a partition key, which the restriction (SPARK-59899) leaves empty;
    • a report that starts with a sort order over a partition transform and has no sort order on a partition key that is not a transform, which SPARK-59995 leaves empty.
  • The conf doc, sql-performance-tuning.md and the 4.4 migration note name the pruned and the transform cases too. An empty report already counts as "no explicit ordering" there.
  • A comment in V2ScanPartitioningAndOrdering said that only None lets the derivation run. It is reworded.
  • A comment in GroupPartitionsExec.outputOrdering said the scan prepends a key-derived ordering. It derives one only when nothing is left of the report. The comment is fixed.
  • DataSourceV2Suite's "ordering and partitioning reporting" now turns the derivation off. Its source breaks the KeyGroupedPartitioning contract, e.g. [1, 1, 3] under the key 1. The test checks the reported ordering and partitioning. Its case with partitioning on i and no ordering expects a sort under groupBy(i), which the derived ordering removes.

The derived ordering leaves partition transforms out, as SPARK-59995 made it do. The SPARK-59995 test in KeyGroupedPartitioningSuite now covers the transform case too.

Why are the changes needed?

With spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled on, the default since SPARK-59396, a keyed scan reports an ordering over its partition key expressions (SPARK-56241). The conf's doc covers a source that "reports a KeyedPartitioning but does not report explicit ordering via SupportsReportOrdering". A scan that implements SupportsReportOrdering and returns an empty array lost that ordering. So did one whose report starts with a pruned column and has no sort order on a partition key. So does one whose report starts with a sort order over a partition transform and has no sort order on a partition key that is not a transform, since SPARK-59995 drops those. A sort-merge join on the partition key then kept a sort above the scan. The results were correct.

Does this PR introduce any user-facing change?

Yes. Such a scan now reports the ordering over its partition keys, so a plan can drop a sort above it. With spark.sql.execution.replaceHashWithSortAgg on, an aggregate over the keys can also be planned as a sort aggregate.

The derivation relies on the KeyGroupedPartitioning contract that every row of a partition has the same partition value. A source that breaks it was already exposed to a wrong derived ordering when it does not implement SupportsReportOrdering. With this change, it is also exposed when it reports an empty ordering, or one that keeps no sort order over the scan output.

How was this patch tested?

New test in KeyGroupedPartitioningSuite, "SPARK-59981: a reported ordering that keeps no sort order falls back to the keys". It joins a scan that reports identity(id) with another identity(id) table. The scan reports three orderings in turn: an empty one, one on s, which the join prunes, and [s, data], where data is not a partition key:

  • with the conf on, the scan keeps the report as it is, its output ordering is id ASC, and the join has no sort above it;
  • with the conf off, there is no ordering and one sort.

It fails on master.

Also ran the KeyGroupedPartitioning* suites, DataSourceV2Suite, the plan merging suites, ProjectedOrderingAndPartitioningSuite, GroupPartitionsExecSuite, PlannerSuite, WriteDistributionAndOrderingSuite, EnsureRequirementsSuite, the plan stability suites, the *JoinSuite suites and AdaptiveQueryExecSuite, 2113 tests in all on the rebased branch, plus dev/lint-scala.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Code (Claude Opus 5.5)

@peter-toth

Copy link
Copy Markdown
Contributor Author

This waits for #59251 (SPARK-59995). The ordering this PR derives can hold a partition transform, such as years(ts), and the k-way merge of GroupPartitionsExec fails on one without that fix. I'll rebase this PR once #59251 is merged, and mark it ready for review then.

@peter-toth
peter-toth force-pushed the SPARK-59981-empty-reported-ordering branch from 482f211 to dd545ce Compare October 8, 2026 19:08
@peter-toth
peter-toth marked this pull request as ready for review October 8, 2026 19:08
@peter-toth

Copy link
Copy Markdown
Contributor Author

Rebased onto master now that #59251 is in. The derived ordering now leaves partition transforms out, and a report that keeps only sort orders over a transform falls back to the keys too. I updated the description.

@HyukjinKwon HyukjinKwon left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Review summary

The outputOrdering rewrite computes the restricted report exactly as the old Some branch did. It then falls back to the same key ordering the None branch already produced, so downstream consumers need no change. The fallback is valid under the documented KeyGroupedPartitioning contract, and the PR description explicitly accepts the extra exposure for contract-breaking sources. The new test asserts that the report reaches the scan unchanged, checks both conf values, and would fail on master.

One P3 documentation gap: the 4.4 migration note for the partitionKeyOrdering default was not extended to the pruned and transform cases that the conf doc and the performance-tuning page now list. A one-sentence update there would keep the three descriptions aligned.

Findings

1 total: 0 P0, 0 P1, 0 P2, 1 P3.

Nit (P3)

  • 4.4 migration note still limits the derived key ordering to scans with no explicit ordering — General
    The PR now derives the key ordering for a reported ordering that keeps no sort order, and updates the conf doc in SQLConf and the docs/sql-performance-tuning.md row to say so. The 4.4 entry in docs/sql-migration-guide.md ("Upgrading from Spark SQL 4.3 to 4.4", line 45) still describes the default-on behavior only for "A V2 scan that reports a keyed partitioning but no explicit ordering".

    A source like the t4 table in the updated SPARK-59995 test reports [years(ts) ASC] under keys (years(ts), id). On branch-4.x (still 4.4.0-SNAPSHOT) it now reports [id ASC], so a Sort above it can disappear and a GROUP BY id can become a sort aggregate. Its owner reading the migration note would not expect that. SPARK-59995 kept this same note in sync when it changed the derivation.

    Suggestion: extend that entry with the same condition you added to the conf doc, for example "...a V2 scan that reports a keyed partitioning but no explicit ordering, or one whose reported ordering starts with a sort order on a column pruned from the scan output or on a partition transform and has no sort order on a non-transform partition key, now also reports itself sorted by its partition key expressions...". Keep the restore instruction as it is. The restore instruction is still correct, so this only affects how users discover the change.

    Verification:

    • Inspection: The migration entry, the SQLConf doc and the performance-tuning row describe the same activation conditions for the derived ordering, and the edited text has no non-ASCII typographic characters.

Verification

  • For empty or fully dropped reports, the fallback emits the same k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending)) ordering that the pre-existing ordering-is-None branch produced, and every non-empty restricted report is returned unchanged.

Generated by Omnigent on Databricks.

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thank you for the fix, @peter-toth. The rewrite keeps the old Some branch's restriction as is and falls back to the same key ordering the None branch already produced, so I think the change is correct under the KeyGroupedPartitioning contract. I left a few inline comments. In summary:

  1. DataSourceV2ScanExecBase.scala L177: Explicit empty report now opts into the derivation

    An empty outputOrdering() used to keep the sort. Now it gets the key ordering by default. A source that breaks the single-key contract, like this suite's own OrderAndPartitionAwareDataSource, can then silently lose a sort it needs, and there is no per-source opt-out. The PR description already accepts this, so the question is whether the migration note should say it explicitly.

  2. DataSourceV2ScanExecBase.scala L176: The derivation depends on column pruning

    With a report [s], the key ordering is derived only when s is pruned. Selecting s, or a merged scan whose output keeps s, brings the sort back. The same happens for a report [data] or [id DESC]. Not blocking. It may be worth a separate JIRA.

  3. DataSourceV2Suite.scala L322: Test coverage

    The conf is turned off for the whole test, although only the (Some("i"), None) case changes.

  4. KeyGroupedPartitioningSuite.scala L4088: Test coverage

    The new test covers only the fallback, not its boundaries ([data], [s, id], the coalescing path).

  5. GroupPartitionsExec.scala L381: K-way merge on a constant key

    With preserveOrderingOnCoalesce on, the newly derived key-only ordering also turns on the k-way merge for these scans, even though that merge cannot gain any ordering.

  6. DataSourceV2ScanExecBase.scala L164: Nit

    ordering.map { ... }.getOrElse(Seq.empty) could be ordering.getOrElse(Nil) with one match.

outputPartitioning match {
case k: KeyedPartitioning
if reported.isEmpty && conf.v2BucketingPartitionKeyOrderingEnabled =>
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.

Before this PR, a source that implements SupportsReportOrdering and returns an empty array (Some(Nil)) got no ordering, so Spark kept the sort above the scan. Now it gets [key ASC] by default, and so does a report that keeps no sort order.

This is fine under the KeyGroupedPartitioning contract. But this suite's own OrderAndPartitionAwareDataSource breaks it: a split holds i = [1, 1, 3] under key 1, and with partitionKeys = j a split holds j = [6, 1, 2] under key 2. That is why DataSourceV2Suite had to turn the conf off. For such a source, a sort-merge join, a sort aggregate (replaceHashWithSortAgg), or flatMapGroups now reads unsorted input and can give wrong results, and the only workaround is the global conf.

I see the PR description accepts this. Could the 4.4 migration note say it explicitly, for example that a source reporting an empty ordering is now also trusted to hold a single key per split? Then a connector owner who relied on an empty report to keep the sort knows to set spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled to false.

}.getOrElse(Seq.empty)
outputPartitioning match {
case k: KeyedPartitioning
if reported.isEmpty && 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.

Not blocking, but I think this gate makes the plan depend on column pruning. With key identity(id) and a report [s]:

  • SELECT t1.data, t2.data ... JOIN ON t1.id = t2.id prunes s, so the report is dropped and the scan derives [id ASC]. No sort.
  • SELECT t1.s, t2.data ... keeps s, so the scan reports [s] and the join adds a sort on id again.
  • A scan merged by PlanMerger whose widened output keeps s loses [id ASC] too. mergeDegradesReporting compares only the reported ordering ([s] on both sides), so it does not see this.

A report [data] or [id DESC] behaves the same way: any surviving sort order suppresses the key ordering, even though id is constant in each split.

Prepending the key orders (keyOrders ++ reported, deduped), as the old GroupPartitionsExec comment described ("key-derived ordering ... already prepended"), would fix these cases. But since orderingSatisfies is prefix-based, it would stop [data] from satisfying a required [data]. So the root fix is probably to make the satisfaction check aware that the key expressions are constant within a split. That seems out of scope here. Should we file a separate JIRA for it, and maybe mention the pruning dependency in the Scaladoc?

// So the ordering derived from the keys is turned off. This test checks the reported one.
withSQLConf(
SQLConf.V2_BUCKETING_ENABLED.key -> "true",
SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "false") {

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.

Only the (Some("i"), None) case changes with this PR. The other cases, [i], [j], [i,j] for both the Scala and the Java source, now run only with the derivation off, so this suite no longer covers how a non-empty report interacts with the default conf.

Could we scope the override to the affected case instead? For example, keep the outer withSQLConf as before and wrap only the (Some("i"), None) case in withSQLConf(SQLConf.V2_BUCKETING_PARTITION_KEY_ORDERING_ENABLED.key -> "false"). Making the source keep a single key per split would also work, but it changes the expected answers.

// one that starts with `s` and goes on with `data`, which is not a partition key.
val sOrder = sort(FieldReference("s"), SortDirection.ASCENDING)
val dataOrder = sort(FieldReference("data"), SortDirection.ASCENDING)
Seq(Seq.empty, Seq(sOrder), Seq(sOrder, dataOrder)).foreach { ordering =>

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.

The three cases here all reduce to nothing, so the test pins the fallback but not its boundaries. Could we add a few cases?

  • [data]: data is in the output and not a key, so the report is kept as [data] and nothing is derived.
  • [s, id]: id survives past the pruned s through rest.filter, so the result is [id] from the report, not from the fallback. Its direction stays as reported, e.g. [s, id DESC] gives [id DESC].
  • A Some(Nil) scan with two splits under one key, under a coalescing GroupPartitionsExec, keeps [id] through preserveKeyOrderingOnCoalesce.

Without them, a change to the reported.isEmpty gate or to the rest.filter branch could change these orderings without failing any test.

// within-partition ordering is fully preserved (including any key-derived ordering that
// `DataSourceV2ScanExecBase` already prepended).
// within-partition ordering is fully preserved, including an ordering that
// `DataSourceV2ScanExecBase` derives from the partition keys.

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.

kWayMergeIsFeasible (L258) is child.outputOrdering.nonEmpty && childIsSafeForKWayMerge. Before this PR only a None report made the derived key-only ordering reach it. Now Some(Nil) and fully dropped reports do too.

With spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled on, tryEnableSortedMerge then plans a SortedMergeCoalescedRDD for such a scan. It merges on a key that is equal in every input of the group, so it gains no ordering beyond what preserveKeyOrderingOnCoalesce already keeps with plain concatenation. It still pays the priority-queue comparisons and turns off columnar output (supportsColumnar is false when usesSortedMerge).

The conf is off by default, so this is minor. Should the merge be feasible only when the child's ordering has a sort order beyond the partition keys?

case _ => prefix
}
case (_, k: KeyedPartitioning) if conf.v2BucketingPartitionKeyOrderingEnabled =>
val reported = ordering.map { o =>

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: spanning Nil already yields an empty prefix and rest, so None needs no separate path, and the two outputPartitioning matches can be one:

val (prefix, rest) = ordering.getOrElse(Nil)
  .span(order => order.references.subsetOf(outputSet) && !holdsTransform(order.child))
outputPartitioning match {
  case k: KeyedPartitioning =>
    val keyExprs = ExpressionSet(k.expressions.filterNot(holdsTransform))
    val kept = prefix ++ rest.filter(order => keyExprs.contains(order.child))
    if (kept.isEmpty && conf.v2BucketingPartitionKeyOrderingEnabled) {
      k.expressions.filterNot(holdsTransform).map(SortOrder(_, Ascending))
    } else {
      kept
    }
  case _ => prefix
}

This also makes it clear why V2ScanPartitioningAndOrdering no longer has to choose between None and Some(Nil).

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants