This is an automated email from the ASF dual-hosted git repository.
voonhous pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 29cdee910a0c [MINOR] Do not return a negative Spark partition when a
hash is Integer.MIN_VALUE (#19776)
29cdee910a0c is described below
commit 29cdee910a0c6c93780c51386c53424b626fbc51
Author: ZIHAN DAI <[email protected]>
AuthorDate: Mon Sep 21 15:38:16 2026 +1000
[MINOR] Do not return a negative Spark partition when a hash is
Integer.MIN_VALUE (#19776)
* [MINOR] Do not return a negative Spark partition when a hash is
Integer.MIN_VALUE
CoalescingPartitioner and PartitionPathRDDPartitioner both derive the
partition
as Math.abs(hash) % numPartitions. Math.abs leaves Integer.MIN_VALUE
negative,
so the expression is negative whenever numPartitions does not divide 2^31,
and
Partitioner#getPartition has to answer inside [0, numPartitions).
Use Math.floorMod, which agrees with the old expression for every hash the
old
one handled correctly. BucketIndexUtil and JavaUpsertPartitioner already
avoid
Math.abs the same way.
* Address review: pin routing against Spark's HashPartitioner
Asserting only that the index is in range passes for abs-mod too. Spark's
HashPartitioner routes with Utils.nonNegativeMod, which is floorMod, so
comparing against it pins the exact index and is a real oracle.
The range-only assertions also let a wrong claim into the PR description:
abs-mod and floorMod differ for roughly half of all negative hashes at every
parallelism, not only for Integer.MIN_VALUE.
---------
Co-authored-by: Zihan Dai <[email protected]>
---
.../apache/hudi/client/CoalescingPartitioner.java | 4 +-
.../bulkinsert/PartitionPathRDDPartitioner.java | 4 +-
.../hudi/client/TestCoalescingPartitioner.java | 38 +++++++++++++++
.../TestPartitionPathRDDPartitioner.java | 54 ++++++++++++++++++++++
4 files changed, 98 insertions(+), 2 deletions(-)
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/CoalescingPartitioner.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/CoalescingPartitioner.java
index 6a449ca39c95..87bf044b146f 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/CoalescingPartitioner.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/CoalescingPartitioner.java
@@ -41,7 +41,9 @@ public class CoalescingPartitioner extends Partitioner {
if (numPartitions == 1) {
return 0;
} else {
- return Math.abs(key.hashCode()) % numPartitions;
+ // Math.abs leaves Integer.MIN_VALUE negative, and a Partitioner must
answer in
+ // [0, numPartitions). floorMod is non-negative for every input.
+ return Math.floorMod(key.hashCode(), numPartitions);
}
}
}
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/execution/bulkinsert/PartitionPathRDDPartitioner.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/execution/bulkinsert/PartitionPathRDDPartitioner.java
index eb835b38c349..a00934e344c0 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/execution/bulkinsert/PartitionPathRDDPartitioner.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/execution/bulkinsert/PartitionPathRDDPartitioner.java
@@ -47,6 +47,8 @@ public class PartitionPathRDDPartitioner extends Partitioner
implements Serializ
@SuppressWarnings("unchecked")
@Override
public int getPartition(Object o) {
- return Math.abs(Objects.hash(partitionPathExtractor.apply(o))) %
numPartitions;
+ // Math.abs leaves Integer.MIN_VALUE negative, and a Partitioner must
answer in
+ // [0, numPartitions). floorMod is non-negative for every input.
+ return Math.floorMod(Objects.hash(partitionPathExtractor.apply(o)),
numPartitions);
}
}
diff --git
a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestCoalescingPartitioner.java
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestCoalescingPartitioner.java
index 9f99742e0c0e..a6bf4e08176d 100644
---
a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestCoalescingPartitioner.java
+++
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestCoalescingPartitioner.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.util.collection.Pair;
import org.apache.hudi.data.HoodieJavaRDD;
import org.apache.hudi.testutils.HoodieClientTestBase;
+import org.apache.spark.HashPartitioner;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.function.FlatMapFunction;
import org.apache.spark.api.java.function.Function;
@@ -184,4 +185,41 @@ public class TestCoalescingPartitioner extends
HoodieClientTestBase {
return booleanIntegerTuple2._2;
}
}
+
+ /**
+ * Spark's own HashPartitioner routes with Utils.nonNegativeMod, which is
floorMod. Pinning the
+ * exact index against it is a real oracle: asserting only that the index is
in range would pass
+ * for abs-mod too, and abs-mod differs from floorMod for roughly half of
all negative hashes.
+ */
+ @Test
+ public void testPartitionMatchesSparkHashPartitioner() {
+ String minValueHashKey = "polygenelubricants";
+ assertEquals(Integer.MIN_VALUE, minValueHashKey.hashCode());
+ // Integer.MIN_VALUE % 2^k == 0, so only 3, 5, 6 and 7 fail on the old
Math.abs expression;
+ // trimming this list to powers of two would stop it exercising the fix.
+ for (int numPartitions : new int[] {1, 2, 3, 4, 5, 6, 7, 8, 16}) {
+ HashPartitioner oracle = new HashPartitioner(numPartitions);
+ for (Object key : new Object[] {minValueHashKey, -1, -2, -3, -5, -100,
0, 1, 100}) {
+ int partition = new
CoalescingPartitioner(numPartitions).getPartition(key);
+ assertTrue(partition >= 0 && partition < numPartitions,
+ "partition " + partition + " out of range for numPartitions " +
numPartitions);
+ assertEquals(oracle.getPartition(key), partition,
+ "key " + key + " at numPartitions " + numPartitions);
+ }
+ }
+ }
+
+ /**
+ * The assertions above only call getPartition. This drives a real shuffle
so the failure the old
+ * expression produced is visible end to end: BypassMergeSortShuffleWriter
indexes its
+ * partitionWriters array with whatever getPartition answers, unguarded.
+ */
+ @Test
+ public void testShuffleSucceedsForMinValueHashKey() {
+ JavaRDD<String> keys =
jsc.parallelize(Collections.singletonList("polygenelubricants"), 1);
+ assertEquals(1, keys.mapToPair(key -> new Tuple2<>(key, key))
+ .partitionBy(new CoalescingPartitioner(3))
+ .count());
+ }
+
}
diff --git
a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/execution/bulkinsert/TestPartitionPathRDDPartitioner.java
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/execution/bulkinsert/TestPartitionPathRDDPartitioner.java
new file mode 100644
index 000000000000..6c3e2a3f32be
--- /dev/null
+++
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/execution/bulkinsert/TestPartitionPathRDDPartitioner.java
@@ -0,0 +1,54 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.execution.bulkinsert;
+
+import org.apache.spark.HashPartitioner;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import java.util.Objects;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class TestPartitionPathRDDPartitioner {
+
+ /**
+ * Objects.hash(x) is 31 + x.hashCode(), so this partition path overflows it
to
+ * Integer.MIN_VALUE. The partitioner hashes that sum rather than the
string, so the oracle is
+ * fed the boxed int. Spark's HashPartitioner routes with
Utils.nonNegativeMod, i.e. floorMod.
+ */
+ // Integer.MIN_VALUE % 2^k == 0, so only 3, 5, 6 and 7 fail on the old
Math.abs expression;
+ // trimming this list to powers of two would stop it exercising the fix.
+ @ParameterizedTest
+ @ValueSource(ints = {1, 2, 3, 4, 5, 6, 7, 8, 16})
+ void testPartitionMatchesSparkHashPartitioner(int numPartitions) {
+ String minValueHashPath = "xfjfxsf";
+ assertEquals(Integer.MIN_VALUE, Objects.hash(minValueHashPath));
+
+ PartitionPathRDDPartitioner partitioner =
+ new PartitionPathRDDPartitioner(o -> minValueHashPath, numPartitions);
+ int partition = partitioner.getPartition(new Object());
+ assertTrue(partition >= 0 && partition < numPartitions,
+ "partition " + partition + " out of range for numPartitions " +
numPartitions);
+ assertEquals(new
HashPartitioner(numPartitions).getPartition(Objects.hash(minValueHashPath)),
+ partition, "numPartitions " + numPartitions);
+ }
+
+}