This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 6d9667af63 [Improve][API] Use HashUtils.bucketIndex in 
MultiTableSinkWriter routing (#12290)
6d9667af63 is described below

commit 6d9667af63e226d7cec44095128e83d52f62299d
Author: Sepuri Sai Krishna <[email protected]>
AuthorDate: Fri Oct 2 08:32:19 2026 +0000

    [Improve][API] Use HashUtils.bucketIndex in MultiTableSinkWriter routing 
(#12290)
---
 .../api/sink/multitablesink/MultiTableSinkWriter.java        | 12 +++++-------
 1 file changed, 5 insertions(+), 7 deletions(-)

diff --git 
a/seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/multitablesink/MultiTableSinkWriter.java
 
b/seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/multitablesink/MultiTableSinkWriter.java
index e17eca1503..590c35bf9a 100644
--- 
a/seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/multitablesink/MultiTableSinkWriter.java
+++ 
b/seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/multitablesink/MultiTableSinkWriter.java
@@ -30,6 +30,7 @@ import 
org.apache.seatunnel.api.table.schema.event.SchemaChangeEvent;
 import org.apache.seatunnel.api.table.type.SeaTunnelRow;
 import org.apache.seatunnel.api.tracing.MDCTracer;
 import org.apache.seatunnel.common.constants.JobMode;
+import org.apache.seatunnel.common.utils.HashUtils;
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -563,9 +564,9 @@ public class MultiTableSinkWriter
      *
      * <ul>
      *   <li>If the table's primary key information is present and the primary 
key field value is
-     *       non-null, the row is routed by {@code (primaryKeyValue.hashCode() 
& Integer.MAX_VALUE)
-     *       % queueSize}, guaranteeing that rows with the same primary key 
always go to the same
-     *       queue for ordered delivery.
+     *       non-null, the row is routed by {@link HashUtils#bucketIndex(int, 
int)} over {@code
+     *       primaryKeyValue.hashCode()} and the queue count, guaranteeing 
that rows with the same
+     *       primary key always go to the same queue for ordered delivery.
      *   <li>If the table's primary key information is present but the actual 
field value is {@code
      *       null}, the row is routed to queue 0.
      *   <li>If the table has no primary key or this is a single-table sink, 
the row is sent to a
@@ -616,10 +617,7 @@ public class MultiTableSinkWriter
             Object object = element.getField(primaryKey.get());
             int index = 0;
             if (object != null) {
-                // Clear the sign bit rather than using Math.abs: 
Math.abs(Integer.MIN_VALUE) is
-                // still Integer.MIN_VALUE, which would yield a negative queue 
index whenever the
-                // queue count is not a power of two.
-                index = (object.hashCode() & Integer.MAX_VALUE) % 
blockingQueues.size();
+                index = HashUtils.bucketIndex(object.hashCode(), 
blockingQueues.size());
             }
             offerRowElement(index, element);
         }

Reply via email to