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

Reply via email to