Skip to content

[SPARK-59524][CORE] Enforce the array-key / HashPartitioner guard in subtractByKey and the Java fakeClassTag key path - #59095

Open
zahed1994 wants to merge 2 commits into
apache:masterfrom
zahed1994:SPARK-59524-array-key-hash-partitioner
Open

zahed1994 wants to merge 2 commits into
apache:masterfrom
zahed1994:SPARK-59524-array-key-hash-partitioner

Conversation

@zahed1994

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

In SPARK-58758 (PR #58758), array keys were prohibited in HashPartitioner and operations relying on array hashing (partitionBy, combineByKeyWithClassTag, aggregateByKey, foldByKey, cogroup). However, two gaps remained:

  1. subtractByKey: PairRDDFunctions.subtractByKey creates a SubtractedRDD directly and defaults to HashPartitioner, but never called failOnHashPartitionerWithArrayKey(p). Because SubtractedRDD matches array keys by reference identity, arrPairs.subtractByKey(arrPairs) silently returned every element instead of subtracting matching keys. This PR adds failOnHashPartitionerWithArrayKey(p) in subtractByKey.
  2. Java API fakeClassTag: In the Java API, JavaSparkContext.parallelizePairs and JavaPairRDD.fromJavaRDD defaulted to fakeClassTag (ClassTag.AnyRef, whose runtime class is Object). Consequently, keyClass.isArray evaluated to false, and byte[] keys bypassed the array-key guard in JavaPairRDD.partitionBy, join, cogroup, and subtractByKey.
    • In JavaSparkContext.parallelizePairs, the key/value ClassTags are now dynamically resolved from non-null elements in the driver-side in-memory collection.
    • In JavaPairRDD, overloaded fromJavaRDD methods accepting keyClass: 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?

  • Array hashing in Java produces reference-based hash codes rather than content-based hash codes, causing incorrect partitioning, grouping, and subtraction.
  • Without these fixes, array keys in subtractByKey silently 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 subtractByKey or Java pair operations (partitionBy, join, cogroup, subtractByKey) on an RDD with array keys and a HashPartitioner will now fail fast with SparkException: [UNSUPPORTED_ARRAY_KEY.HASH_PARTITIONER] rather than silently returning incorrect results.

How was this patch tested?

  • Added subtractByKey test cases in core/src/test/scala/org/apache/spark/PartitioningSuite.scala.
  • Added testArrayKeyUnderHashPartitionerFails in core/src/test/java/test/org/apache/spark/JavaAPISuite.java.
  • Ran:
    • build/sbt "core/testOnly org.apache.spark.PartitioningSuite"
    • build/sbt "core/testOnly test.org.apache.spark.JavaAPISuite"
    • build/sbt core/scalastyle

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

No.

…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
@zahed1994

zahed1994 commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor Author

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):

  1. Adds failOnHashPartitionerWithArrayKey(p) in subtractByKey.
  2. Resolves actual key/value ClassTags in JavaSparkContext.parallelizePairs from the in-memory driver collection.
  3. Adds @Since("4.4.0") typed fromJavaRDD(rdd, keyClass) overloads in JavaPairRDD to handle JVM type erasure.

Thank you!

@uros-b uros-b 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.

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))
}

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.

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.

@zahed1994

Copy link
Copy Markdown
Contributor Author

Thanks for the sharp catch, @uros-b - that is a subtle and critical edge case.

I have updated the PR accordingly:

  1. Value side (ctagV): Reverted unconditionally to fakeClassTag. Because HashPartitioner is fundamentally a key-only guard, inferring value types was unnecessary and risked ArrayStoreException when collecting heterogeneous values.
  2. Key side (ctagK): Gated the specialization strictly to array types (isArray). For any non-array or heterogeneous keys (e.g. Number with Integer and Double), it falls back to fakeClassTag, ensuring keys().collect() continues allocating Object[] as before without ArrayStoreException.
  3. Tests: Added a regression test in JavaAPISuite covering heterogeneous keys and values to ensure neither throws ArrayStoreException, while array keys under HashPartitioner continue to be blocked.

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.

2 participants