[SPARK-59524][CORE] Enforce the array-key / HashPartitioner guard in subtractByKey and the Java fakeClassTag key path - #59095
Conversation
…subtractByKey and the Java fakeClassTag key path ### What changes were proposed in this pull request? This PR addresses two gaps in array-key prohibition under HashPartitioner: 1. Adds failOnHashPartitionerWithArrayKey(p) in PairRDDFunctions.subtractByKey before constructing SubtractedRDD. 2. In JavaSparkContext.parallelizePairs, dynamically derives ClassTag from non-null keys/values in the in-memory driver collection instead of unconditionally using fakeClassTag. 3. In JavaPairRDD, adds fromJavaRDD overloads accepting Class[K] and Class[V] to allow Java callers to provide explicit key type information across Java generic erasure. ### Why are the changes needed? In SPARK-58758, subtractByKey was missed when guarding PairRDDFunctions methods with failOnHashPartitionerWithArrayKey. Furthermore, Java pair RDDs using fakeClassTag had keyClass evaluated as Object, bypassing the array-key guard on byte[] keys. ### Does this PR introduce any user-facing change? Yes, using array keys under HashPartitioner in subtractByKey and Java pair RDD operations now fails fast with UNSUPPORTED_ARRAY_KEY.HASH_PARTITIONER instead of producing silent incorrect results. ### How was this patch tested? - core/testOnly org.apache.spark.PartitioningSuite - core/testOnly test.org.apache.spark.JavaAPISuite
|
Hi @dongjoon-hyun and @uros-b, Friendly ping when you have a moment. This PR addresses the follow-up items discussed in #58758 (SPARK-59524):
Thank you! |
uros-b
left a comment
There was a problem hiding this comment.
This is a correct, well-scoped follow-up that closes the two gaps identified on the parent PR; guard placement, @SInCE version, and error/test wiring are all right. Before merge, we should weigh one edge: resolving the value-side ClassTag from the first element is unnecessary for the key-only guard and needlessly regresses heterogeneous collections to ArrayStoreException on keys()/values().collect() - prefer leaving ctagV = fakeClassTag.
| case None => fakeClassTag[V] | ||
| } | ||
| JavaPairRDD.fromRDD(sc.parallelize(seq, numSlices)) | ||
| } |
There was a problem hiding this comment.
Resolving ClassTag[K] AND ClassTag[V] from the first non-null element's runtime class (vs the prior fakeClassTag -> Object) stores a leaf ClassTag. For a heterogeneous collection (e.g. List<Tuple2<Number, V>> with keys [Integer(1), Double(2.0)], legal Java) a later keys().collect() / values().collect() allocates Array[K]/Array[V] via that specific tag (rdd.map(_._1) -> iter.toArray) and can throw ArrayStoreException where Object[] previously succeeded - a newly introduced behavior change in a widely used public API, uncovered by the tests.
… keys and all values in parallelizePairs
|
Thanks for the sharp catch, @uros-b - that is a subtle and critical edge case. I have updated the PR accordingly:
|
What changes were proposed in this pull request?
In SPARK-58758 (PR #58758), array keys were prohibited in
HashPartitionerand operations relying on array hashing (partitionBy,combineByKeyWithClassTag,aggregateByKey,foldByKey,cogroup). However, two gaps remained:subtractByKey:PairRDDFunctions.subtractByKeycreates aSubtractedRDDdirectly and defaults toHashPartitioner, but never calledfailOnHashPartitionerWithArrayKey(p). BecauseSubtractedRDDmatches array keys by reference identity,arrPairs.subtractByKey(arrPairs)silently returned every element instead of subtracting matching keys. This PR addsfailOnHashPartitionerWithArrayKey(p)insubtractByKey.fakeClassTag: In the Java API,JavaSparkContext.parallelizePairsandJavaPairRDD.fromJavaRDDdefaulted tofakeClassTag(ClassTag.AnyRef, whose runtime class isObject). Consequently,keyClass.isArrayevaluated tofalse, andbyte[]keys bypassed the array-key guard inJavaPairRDD.partitionBy,join,cogroup, andsubtractByKey.JavaSparkContext.parallelizePairs, the key/valueClassTags are now dynamically resolved from non-null elements in the driver-side in-memory collection.JavaPairRDD, overloadedfromJavaRDDmethods acceptingkeyClass: Class[K]/valueClass: Class[V]are added with@Since("4.4.0")to allow Java callers to explicitly provide key type information across JVM generic type erasure.Why are the changes needed?
subtractByKeysilently produce incorrect results, and Java pair RDDs with array keys completely evade the safety guard introduced in SPARK-58758.Does this PR introduce any user-facing change?
Yes. Calling
subtractByKeyor Java pair operations (partitionBy,join,cogroup,subtractByKey) on an RDD with array keys and aHashPartitionerwill now fail fast withSparkException: [UNSUPPORTED_ARRAY_KEY.HASH_PARTITIONER]rather than silently returning incorrect results.How was this patch tested?
subtractByKeytest cases incore/src/test/scala/org/apache/spark/PartitioningSuite.scala.testArrayKeyUnderHashPartitionerFailsincore/src/test/java/test/org/apache/spark/JavaAPISuite.java.build/sbt "core/testOnly org.apache.spark.PartitioningSuite"build/sbt "core/testOnly test.org.apache.spark.JavaAPISuite"build/sbt core/scalastyleWas this patch authored or co-authored using generative AI tooling?
No.