Skip to content

Commit 751a496

Browse files
anuragmantridongjoon-hyun
authored andcommitted
[SPARK-59721][SQL] Ignore reported partitioning and ordering that reference unresolvable columns
### What changes were proposed in this pull request? When a V2 scan reports a `KeyGroupedPartitioning` or an `ordering` that references a column the relation can't resolve, the query used to fail with `Unable to resolve X given [...]`. This PR makes `V2ScanPartitioningAndOrdering` check the columns referenced by the reported partition keys and ordering (recursively, including nested transform arguments) against the relation before converting them. If any can't be resolved, including a missing nested field, an ambiguous name, or an empty name, the rule logs a warning naming the relation, the scan class, and those columns, and ignores the report. The partitioning falls back to `UnknownPartitioning`, and the ordering is dropped. This PR adds a `resolveRefOpt`, and makes `resolveRef` a thin wrapper around it. `resolveRefOpt` still throws for a missing nested field or an ambiguous name, so the rule treats those errors as unresolvable and includes the error condition in the warning. The `toCatalyst*` methods are unchanged. Keys whose columns all resolve take the existing conversion path unchanged, so an unsupported expression type (e.g. `id + 1`) or a nested transform whose function can't be loaded still fails the query as before. Because the whole report is checked before any key is converted, a report that also references an unresolvable column falls back with a warning instead, e.g. `[id + 1, identity(missing)]`. The Javadoc of `KeyGroupedPartitioning#keys()` and `SupportsReportOrdering#outputOrdering()` now says that Spark resolves the column references against the table columns, including columns pruned from the scan output, and ignores the whole report if any of them can't be resolved. The docs of `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` now mention that an ignored ordering can be replaced by one derived from the partition keys. ### Why are the changes needed? Found while reviewing SPARK-58111 ([PR #55518](#55518)). This is a general, pre-existing bug in how `V2ScanPartitioningAndOrdering` handles reported partitioning and ordering, not something introduced by #55518. ### Does this PR introduce _any_ user-facing change? Yes. A query against a data source whose scan reports a `KeyGroupedPartitioning` with a key that can't be resolved previously failed outright. It now succeeds, with partitioning degraded to `UnknownPartitioning` instead, and Spark logs a warning naming the relation, the scan class, and the columns that can't be resolved. The same applies to an ambiguous or empty column name. A report that combines an unsupported expression with an unresolvable column, e.g. `missing + 1` or `[id + 1, identity(missing)]`, used to fail with `_LEGACY_ERROR_TEMP_3054` and now falls back with a warning. An unsupported expression over resolvable columns still fails as before. If the scan keeps its partitioning but its reported ordering is ignored, Spark derives the ordering from the partition keys when `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` is on (the default). ### How was this patch tested? New tests in `V2ExpressionUtilsSuite`: - `resolveRefOpt` returns `None` for an unresolvable reference. - `resolveRefOpt` throws `FIELD_NOT_FOUND` for a missing nested field and `AMBIGUOUS_REFERENCE` for an ambiguous reference. New tests in `KeyGroupedPartitioningSuite`: - A join against a table whose scan reports unresolvable partition keys succeeds, plans a shuffle instead of a storage-partitioned join, and logs a warning naming the relation, the scan class, and only the unresolvable columns. The cases are: - a `bucket` key, a nested transform key, and a multi-key partitioning with one bad key; - duplicate references, and two distinct unresolvable columns; - a missing nested field; - an unsupported expression next to or over an unresolvable column (`[id + 1, identity(missing)]`, `missing + 1`); - a connector-defined transform that hides the reference from `children()`. Resolvable keys (`id`, case-insensitive `ID`, and a connector-defined transform with an extra reference only in `children()`) keep the partitioning and plan a storage-partitioned join without a shuffle. - A reported ordering that can't be resolved is ignored as a whole with a warning, and the scan still derives `id ASC` from its kept partitioning. The cases are a column that doesn't exist, one hidden from a connector-defined sort order's `children()`, an empty name, and a two-key ordering with one unresolvable key (`[id, missing]`). Resolvable orderings (`id`, case-insensitive `ID`, the metadata column `index`, and a connector-defined sort order with an extra reference only in `children()`) are kept without a warning. - A reported partition key of an unsupported expression type (`id + 1`), or a nested transform whose function can't be loaded (`f(g(id))`), still fails the query with `_LEGACY_ERROR_TEMP_3054`. New test in `DataSourceV2Suite`: - An unresolvable ordering reported through connector-defined `SortOrder`, `Transform`, and `NamedReference` classes (`OrderAndPartitionAwareDataSource` and its Java counterpart) is ignored with a warning, and the scan relation keeps no ordering. Also ran `V2ExpressionUtilsSuite`, `KeyGroupedPartitioningSuite`, `WriteDistributionAndOrderingSuite`, `DataSourceV2Suite`, `MergeSubplansSuite` (comment update only), `SQLConfSuite`, and `ProjectedOrderingAndPartitioningSuite` to confirm no regression in the partitioning/SPJ, ordering, and write paths this change touches. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code (Sonnet 5, Opus 5.5) Verified manually by me. Closes #58979 from anuragmantri/v2expr-opt-resolve-degrade. Authored-by: Anurag Mantripragada <amantripragada@apple.com> Signed-off-by: Dongjoon Hyun <dongjoon@apache.org> (cherry picked from commit 9de605f) Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
1 parent d8d72d2 commit 751a496

12 files changed

Lines changed: 409 additions & 44 deletions

File tree

‎docs/sql-performance-tuning.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -708,7 +708,7 @@ The following SQL properties enable Storage Partition Join in different join que
708708
<td><code>spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled</code></td>
709709
<td>true</td>
710710
<td>
711-
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. 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.
711+
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.
712712
</td>
713713
<td>4.2.0</td>
714714
</tr>

‎sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/SupportsReportOrdering.java‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,10 @@ public interface SupportsReportOrdering extends Scan {
3535

3636
/**
3737
* Returns the order in each partition of this data source scan.
38+
* <p>
39+
* Spark resolves the column references in these sort orders against the table columns,
40+
* including columns pruned from the scan output. If any of them cannot be resolved, Spark
41+
* ignores the whole reported ordering and logs a warning.
3842
*/
3943
SortOrder[] outputOrdering();
4044
}

‎sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/partitioning/KeyGroupedPartitioning.java‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,10 @@ public KeyGroupedPartitioning(Expression[] keys, int numPartitions) {
4343

4444
/**
4545
* Returns the partition transform expressions for this partitioning.
46+
* <p>
47+
* Spark resolves the column references in these expressions against the table columns,
48+
* including columns pruned from the scan output. If any of them cannot be resolved, Spark
49+
* ignores the whole reported partitioning and logs a warning.
4650
*/
4751
public Expression[] keys() {
4852
return keys;

‎sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/V2ExpressionUtils.scala‎

Lines changed: 15 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -46,15 +46,22 @@ import org.apache.spark.util.ArrayImplicits._
4646
object V2ExpressionUtils extends SQLConfHelper with Logging {
4747
import org.apache.spark.sql.connector.catalog.CatalogV2Implicits.MultipartIdentifierHelper
4848

49+
/**
50+
* Variant of `resolveRef` that returns `None` if no attribute matches the reference. It still
51+
* throws for a nested-field extraction error, such as a missing nested field, or an ambiguous
52+
* reference.
53+
*/
54+
private[sql] def resolveRefOpt(
55+
ref: NamedReference, plan: LogicalPlan): Option[NamedExpression] = {
56+
plan.resolve(ref.fieldNames.toImmutableArraySeq, conf.resolver)
57+
}
58+
4959
def resolveRef[T <: NamedExpression](ref: NamedReference, plan: LogicalPlan): T = {
50-
plan.resolve(ref.fieldNames.toImmutableArraySeq, conf.resolver) match {
51-
case Some(namedExpr) =>
52-
namedExpr.asInstanceOf[T]
53-
case None =>
54-
val name = ref.fieldNames.toImmutableArraySeq.quoted
55-
val outputString = plan.output.map(_.name).mkString(",")
56-
throw QueryCompilationErrors.cannotResolveAttributeError(name, outputString)
57-
}
60+
resolveRefOpt(ref, plan).getOrElse {
61+
val name = ref.fieldNames.toImmutableArraySeq.quoted
62+
val outputString = plan.output.map(_.name).mkString(",")
63+
throw QueryCompilationErrors.cannotResolveAttributeError(name, outputString)
64+
}.asInstanceOf[T]
5865
}
5966

6067
def resolveRefs[T <: NamedExpression](refs: Seq[NamedReference], plan: LogicalPlan): Seq[T] = {

‎sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2590,7 +2590,8 @@ object SQLConf {
25902590
buildConf("spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled")
25912591
.doc("When enabled, Spark derives output ordering from the partition key expressions of " +
25922592
"a V2 data source that reports a KeyedPartitioning but does not report explicit ordering " +
2593-
"via SupportsReportOrdering. Within a single partition all rows share the same key " +
2593+
"via SupportsReportOrdering, or reports one that Spark ignores because it references a " +
2594+
"column that cannot be resolved. Within a single partition all rows share the same key " +
25942595
s"value, so the data is trivially sorted by those expressions. Requires " +
25952596
s"${V2_BUCKETING_ENABLED.key} to be enabled.")
25962597
.version("4.2.0")

‎sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/V2ExpressionUtilsSuite.scala‎

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ import org.apache.spark.SparkFunSuite
2121
import org.apache.spark.sql.AnalysisException
2222
import org.apache.spark.sql.catalyst.plans.logical.LocalRelation
2323
import org.apache.spark.sql.connector.expressions._
24-
import org.apache.spark.sql.types.StringType
24+
import org.apache.spark.sql.types.{IntegerType, StringType, StructType}
2525

2626
class V2ExpressionUtilsSuite extends SparkFunSuite {
2727

@@ -37,4 +37,30 @@ class V2ExpressionUtilsSuite extends SparkFunSuite {
3737
}
3838
assert(exc.message.contains("v2Fun(a) ASC NULLS FIRST is not currently supported"))
3939
}
40+
41+
test("SPARK-59721: resolveRefOpt returns None for an unresolvable reference") {
42+
val plan = LocalRelation(AttributeReference("a", StringType)())
43+
assert(V2ExpressionUtils.resolveRefOpt(FieldReference("a"), plan).isDefined)
44+
assert(V2ExpressionUtils.resolveRefOpt(FieldReference("missing"), plan).isEmpty)
45+
}
46+
47+
test("SPARK-59721: resolveRefOpt throws for a missing nested field or an ambiguous reference") {
48+
val structPlan =
49+
LocalRelation(AttributeReference("s", new StructType().add("x", IntegerType))())
50+
checkError(
51+
exception = intercept[AnalysisException] {
52+
V2ExpressionUtils.resolveRefOpt(FieldReference("s.missing"), structPlan)
53+
},
54+
condition = "FIELD_NOT_FOUND",
55+
parameters = Map("fieldName" -> "`missing`", "fields" -> "`x`"))
56+
57+
val ambiguousPlan = LocalRelation(
58+
AttributeReference("a", StringType)(), AttributeReference("a", StringType)())
59+
checkError(
60+
exception = intercept[AnalysisException] {
61+
V2ExpressionUtils.resolveRefOpt(FieldReference("a"), ambiguousPlan)
62+
},
63+
condition = "AMBIGUOUS_REFERENCE",
64+
parameters = Map("name" -> "`a`", "referenceNames" -> "[`a`, `a`]"))
65+
}
4066
}

‎sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala‎

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -135,15 +135,15 @@ trait DataSourceV2ScanExecBase
135135

136136
/**
137137
* Returns the output ordering for this scan. When the source reports ordering via
138-
* `SupportsReportOrdering`, that ordering may reference columns pruned out of the scan output
139-
* (see V2ScanPartitioningAndOrdering), while a consumer binds it against the output, e.g. the
140-
* k-way merge of `GroupPartitionsExec`. Ordering is prefix-based, so the leading run of sort
141-
* orders over the output is kept and the rest is dropped, except the sort orders on a partition
142-
* key when the output partitioning is a `KeyedPartitioning`: each partition holds a single key,
143-
* so those still hold. Otherwise, when the output partitioning is a `KeyedPartitioning` and
144-
* `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled` is on, each partition
145-
* contains rows where the key expressions evaluate to a single constant value, so the data
146-
* is trivially sorted by those expressions within the partition.
138+
* `SupportsReportOrdering` and `V2ScanPartitioningAndOrdering` keeps it, that ordering may
139+
* reference columns pruned out of the scan output, while a consumer binds it against the
140+
* output, e.g. the k-way merge of `GroupPartitionsExec`. Ordering is prefix-based, so the
141+
* leading run of sort orders over the output is kept and the rest is dropped, except the sort
142+
* orders on a partition key when the output partitioning is a `KeyedPartitioning`: each
143+
* partition holds a single key, so those still hold. Otherwise, when the output partitioning
144+
* is a `KeyedPartitioning` and `spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled`
145+
* is on, each partition contains rows where the key expressions evaluate to a single constant
146+
* value, so the data is trivially sorted by those expressions within the partition.
147147
*/
148148
override def outputOrdering: Seq[SortOrder] = {
149149
(ordering, outputPartitioning) match {

‎sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/V2ScanPartitioningAndOrdering.scala‎

Lines changed: 69 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -17,11 +17,14 @@
1717
package org.apache.spark.sql.execution.datasources.v2
1818

1919
import org.apache.spark.internal.Logging
20-
import org.apache.spark.internal.LogKeys.CLASS_NAME
20+
import org.apache.spark.internal.LogKeys.{CLASS_NAME, COLUMN_NAMES, RELATION_NAME}
21+
import org.apache.spark.sql.AnalysisException
2122
import org.apache.spark.sql.catalyst.expressions.V2ExpressionUtils
2223
import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
2324
import org.apache.spark.sql.catalyst.rules.Rule
2425
import org.apache.spark.sql.catalyst.trees.TreePattern.DATA_SOURCE_V2_SCAN_RELATION
26+
import org.apache.spark.sql.connector.catalog.CatalogV2Implicits.MultipartIdentifierHelper
27+
import org.apache.spark.sql.connector.expressions.{Expression => V2Expression, NamedReference, SortOrder => V2SortOrder, Transform}
2528
import org.apache.spark.sql.connector.read.{SupportsReportOrdering, SupportsReportPartitioning}
2629
import org.apache.spark.sql.connector.read.partitioning.{KeyGroupedPartitioning, UnknownPartitioning}
2730
import org.apache.spark.util.ArrayImplicits._
@@ -47,9 +50,19 @@ object V2ScanPartitioningAndOrdering extends Rule[LogicalPlan] with Logging {
4750
if d.keyGroupedPartitioning.isEmpty =>
4851
val catalystPartitioning = scan.outputPartitioning() match {
4952
case kgp: KeyGroupedPartitioning =>
50-
val partitioning = sequenceToOption(
51-
kgp.keys().map(V2ExpressionUtils.toCatalystOpt(_, relation, relation.funCatalog))
52-
.toImmutableArraySeq)
53+
val keys = kgp.keys().toImmutableArraySeq
54+
val unresolvedColumns = unresolvableColumns(keys, relation)
55+
val partitioning = if (unresolvedColumns.nonEmpty) {
56+
logWarning(
57+
log"Spark ignores the KeyGroupedPartitioning reported by " +
58+
log"${MDC(RELATION_NAME, relation.name)} (scan " +
59+
log"${MDC(CLASS_NAME, scan.getClass.getName)}) because the partition key columns " +
60+
log"cannot be resolved: ${MDC(COLUMN_NAMES, unresolvedColumns.mkString(", "))}.")
61+
None
62+
} else {
63+
sequenceToOption(
64+
keys.map(V2ExpressionUtils.toCatalystOpt(_, relation, relation.funCatalog)))
65+
}
5366
// Keep the partitioning when at least one of its keys is still in the scan output: the
5467
// scan projects the pruned key positions away when reporting its physical output
5568
// partitioning (see DataSourceV2ScanExecBase.outputPartitioning). When no key survives,
@@ -75,12 +88,57 @@ object V2ScanPartitioningAndOrdering extends Rule[LogicalPlan] with Logging {
7588
private def ordering(plan: LogicalPlan) = plan.transformDownWithPruning(
7689
_.containsPattern(DATA_SOURCE_V2_SCAN_RELATION)) {
7790
case d @ ExtractV2ScanInfo(relation, scan: SupportsReportOrdering, _) =>
78-
// The ordering is kept as reported, even where it references columns pruned out of the scan
79-
// output: truncating it here would also drop the sort orders on a partition key past a
80-
// pruned column, which still hold. `DataSourceV2ScanRelation.doCanonicalize` and
81-
// `DataSourceV2ScanExecBase.outputOrdering` restrict it to the scan output instead.
82-
val ordering =
83-
V2ExpressionUtils.toCatalystOrdering(scan.outputOrdering(), relation, relation.funCatalog)
84-
d.copy(ordering = Some(ordering))
91+
val reportedOrdering = scan.outputOrdering()
92+
val unresolvedColumns = unresolvableColumns(reportedOrdering.toImmutableArraySeq, relation)
93+
if (unresolvedColumns.nonEmpty) {
94+
logWarning(
95+
log"Spark ignores the ordering reported by ${MDC(RELATION_NAME, relation.name)} " +
96+
log"(scan ${MDC(CLASS_NAME, scan.getClass.getName)}) because the ordering columns " +
97+
log"cannot be resolved: ${MDC(COLUMN_NAMES, unresolvedColumns.mkString(", "))}.")
98+
// Use None, not Some(Nil): only None lets DataSourceV2ScanExecBase.outputOrdering derive an
99+
// ordering from a kept key-grouped partitioning, which does not depend on the dropped
100+
// report.
101+
d.copy(ordering = None)
102+
} else {
103+
// The ordering is kept as reported, even where it references columns pruned out of the
104+
// scan output: truncating it here would also drop the sort orders on a partition key past a
105+
// pruned column, which still hold. `DataSourceV2ScanRelation.doCanonicalize` and
106+
// `DataSourceV2ScanExecBase.outputOrdering` restrict it to the scan output instead.
107+
val ordering =
108+
V2ExpressionUtils.toCatalystOrdering(reportedOrdering, relation, relation.funCatalog)
109+
d.copy(ordering = Some(ordering))
110+
}
111+
}
112+
113+
private def unresolvableColumns(
114+
exprs: Seq[V2Expression],
115+
relation: LogicalPlan): Seq[String] = {
116+
// Walk the accessors the conversion reads, `Transform.arguments()` and
117+
// `SortOrder.expression()`, instead of `V2Expression.references()` or `children()`: a connector
118+
// can override those to disagree with the conversion, and `ApplyTransform.references()` returns
119+
// only its top-level arguments, which would miss `f(g(missing))`. Other expressions, such as
120+
// `GeneralScalarExpression`, which the conversion rejects, go through `children()`.
121+
def collectReferences(expr: V2Expression): Seq[NamedReference] = expr match {
122+
case ref: NamedReference => Seq(ref)
123+
case t: Transform => t.arguments().toImmutableArraySeq.flatMap(collectReferences)
124+
case s: V2SortOrder => collectReferences(s.expression())
125+
case other => other.children().toImmutableArraySeq.flatMap(collectReferences)
126+
}
127+
exprs.flatMap(collectReferences).flatMap { ref =>
128+
val parts = ref.fieldNames.toImmutableArraySeq
129+
if (parts.isEmpty) {
130+
// An empty name can be neither quoted nor resolved: both throw.
131+
Some("<empty>")
132+
} else {
133+
val name = parts.quoted
134+
try {
135+
Option.when(V2ExpressionUtils.resolveRefOpt(ref, relation).isEmpty)(name)
136+
} catch {
137+
// A nested-field extraction error or an ambiguous reference throws instead of returning
138+
// None.
139+
case e: AnalysisException => Some(s"$name (${e.getCondition})")
140+
}
141+
}
142+
}.distinct
85143
}
86144
}

‎sql/core/src/main/scala/org/apache/spark/sql/execution/planmerging/PlanMerger.scala‎

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -980,12 +980,13 @@ class PlanMerger(
980980
// Only the reported EXPRESSIONS are compared, which is all a DataSourceV2ScanRelation carries;
981981
// the merged scan's split count and partition values can still differ from an input's, since it
982982
// may push a different best-effort filter and so prune differently. And a report the merged scan
983-
// GAINS is not a degradation either. For partitioning that is because an input dropped its own
984-
// only where none of its keys survived in that scan's output (V2ScanPartitioningAndOrdering's
985-
// partitioning pass keeps the report whenever any key survives), not because the source stopped
986-
// reporting; the ordering pass has no such guard, so an ordering report is
987-
// never dropped by pruning and a gained one can only come from the source. Either way, keeping it
988-
// is exactly the win this merge is after.
983+
// GAINS is not a degradation either. Both V2ScanPartitioningAndOrdering passes drop a report
984+
// that references a column the relation cannot resolve, but that does not depend on pruning.
985+
// Beyond that, an input dropped its partitioning only where none of its keys survived in that
986+
// scan's output (the partitioning pass keeps the report whenever any key survives), not because
987+
// the source stopped reporting; the ordering pass has no such guard, so an ordering report is
988+
// never dropped by pruning and a gained one can only come from the source. Either way, keeping
989+
// it is exactly the win this merge is after.
989990
private def mergeDegradesReporting(
990991
merged: DataSourceV2ScanRelation,
991992
requiredKeyGroupedPartitioning: Seq[Expression],

‎sql/core/src/test/scala/org/apache/spark/sql/connector/DataSourceV2Suite.scala‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import java.util.OptionalLong
2424

2525
import scala.jdk.CollectionConverters._
2626

27+
import org.apache.logging.log4j.Level
2728
import test.org.apache.spark.sql.connector._
2829

2930
import org.apache.spark.{SparkException, SparkUnsupportedOperationException}
@@ -397,6 +398,34 @@ class DataSourceV2Suite extends SharedSparkSession with AdaptiveSparkPlanHelper
397398
}
398399
}
399400

401+
test("SPARK-59721: ignore an unresolvable ordering reported with connector-defined classes") {
402+
Seq(
403+
classOf[OrderAndPartitionAwareDataSource],
404+
classOf[JavaOrderAndPartitionAwareDataSource]
405+
).foreach { cls =>
406+
withClue(cls.getName) {
407+
val df = spark.read
408+
.option("partitionKeys", "i")
409+
.option("orderKeys", "missing")
410+
.format(cls.getName)
411+
.load()
412+
val logAppender = new LogAppender("unresolvable reported ordering")
413+
withLogAppender(logAppender, level = Some(Level.WARN)) {
414+
checkAnswer(df, Seq(Row(1, 4), Row(1, 5), Row(3, 5), Row(2, 6), Row(4, 1), Row(4, 2)))
415+
}
416+
val scanRelation = getScanRelation(df)
417+
val expectedWarning =
418+
s"Spark ignores the ordering reported by ${scanRelation.name} " +
419+
s"(scan ${scanRelation.scan.getClass.getName}) because the ordering columns " +
420+
"cannot be resolved: missing."
421+
val warnings = logAppender.loggingEvents.map(_.getMessage.getFormattedMessage)
422+
.filter(_.contains("columns cannot be resolved"))
423+
assert(warnings.toSet == Set(expectedWarning), warnings)
424+
assert(scanRelation.ordering.isEmpty)
425+
}
426+
}
427+
}
428+
400429
test ("statistics report data source") {
401430
Seq(classOf[ReportStatisticsDataSource], classOf[JavaReportStatisticsDataSource]).foreach {
402431
cls =>

0 commit comments

Comments
 (0)