peter-toth commented on code in PR #58899:
URL: https://github.com/apache/spark/pull/58899#discussion_r4194252359
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/physical/partitioning.scala:
##########
@@ -972,6 +973,23 @@ object KeyedPartitioning {
KeyLayout(comparablePartitionKeys, factory.dataTypes, isGrouped,
isCollapsed = false))
}
+ // The key's arity is implicit in `HasPartitionKey` and read positionally
downstream: a key wider
+ // than `expressions` raises an `ArrayIndexOutOfBoundsException` in the
struct hash unless the
+ // scan reports one partition, and a narrower one silently groups too
loosely until two keys tie
+ // on their leading fields and the ordering reads past its end.
+ def checkPartitionKeyArity(
+ expressions: Seq[Expression],
+ partitionKeys: Seq[InternalRow]): Unit = {
+ partitionKeys.foreach { key =>
+ if (key.numFields != expressions.length) {
+ throw new SparkException("Data source reported a partition key with " +
Review Comment:
**Finding 1.** The description says the new error names "the contract, the
arity it got, and the expressions it was measured against". The message gives
the two counts only. It also doesn't say which scan reported the key, and that
is what a user with several DSv2 sources needs first.
Two small changes would cover both:
- list the expressions, e.g. `expressions.map(_.sql).mkString(", ")`;
- let the two call sites that know the scan pass it in, e.g.
`scan.description()` from `DataSourceV2ScanExecBase.reportedKeyedPartitioning`
and `PushDownUtils.replanWithRuntimeFilters`.
Or trim the description to what the message says.
##########
sql/catalyst/src/main/java/org/apache/spark/sql/connector/read/HasPartitionKey.java:
##########
@@ -36,7 +36,9 @@
* <p>
* It is implementor's responsibility to ensure that when an input partition
implements this
* interface, its records all have the same value for the partition keys.
Spark doesn't check
- * this property.
+ * that the records agree. It does check the key's width: the row returned by
+ * {@link #partitionKey()} must hold one field per partition expression the
scan reports through
+ * {@link SupportsReportPartitioning}, and Spark rejects a key of any other
width.
Review Comment:
**Finding 2.** "Spark rejects a key of any other width" reads as
unconditional. Spark only reads the key when it uses the reported partitioning:
`spark.sql.sources.v2.bucketing.enabled` is on,
`KeyedPartitioning.supportsExpressions` accepts the expressions, and every
split implements `HasPartitionKey`
(`DataSourceV2ScanExecBase.reportedKeyedPartitioning`). The migration guide's
way back, turning bucketing off, relies on exactly that.
```suggestion
* {@link SupportsReportPartitioning}. When Spark uses that partitioning, it
rejects a key of
* any other width.
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]