SEPURI-SAI-KRISHNA opened a new pull request, #11721:
URL: https://github.com/apache/seatunnel/pull/11721

   ### Purpose of this pull request
   
   Closes #11720.
   
   `MultiTableSinkWriter.write` routed rows with `Math.abs(object.hashCode()) % 
blockingQueues.size()`. `Math.abs(Integer.MIN_VALUE)` returns 
`Integer.MIN_VALUE` — still negative, since `-Integer.MIN_VALUE` overflows 
`int` — so any primary key hashing to `Integer.MIN_VALUE` produced a negative 
index and failed the sink writer with `IndexOutOfBoundsException` whenever the 
queue count was not a power of two.
   
   An `INT` primary key of `-2147483648`, a `BIGINT` primary key of 
`Long.MIN_VALUE`, and any string primary key hashing to `Integer.MIN_VALUE` all 
hit it. Because `multi_table_sink_replica` defaults to `1` and powers of two 
are also safe, this only surfaces once that option is tuned to a 
non-power-of-two value — but it affects every multi-table sink, since the 
routing is shared code in `seatunnel-api`.
   
   The fix clears the sign bit instead of negating:
   
   ```java
   index = (object.hashCode() & Integer.MAX_VALUE) % blockingQueues.size();
   ```
   
   This is the same idiom #2921 introduced for 
`FileSourceSplitEnumerator.getSplitOwner`, and it is already used in 12 other 
places (Paimon, Kafka, DynamoDB, HBase, Fluss, Typesense, TiDB CDC, MongoDB, 
Iceberg, JDBC, Easysearch, and the Hazelcast metrics store). 
`MultiTableSinkWriter` was the last site still using the unsafe form.
   
   The class Javadoc, which documented the old formula, has been updated to 
match.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Only in a narrow and unobservable sense, so: no.
   
   Keys with a non-negative hash are entirely unaffected — `h & 
Integer.MAX_VALUE == h` for every non-negative `h`. Keys with a negative hash 
may now land in a different queue than before, since the old expression 
computed `-h` and the new one computes `h + 2^31`. That is not observable in 
the output: the routing contract is only that a given primary key 
deterministically maps to one queue, so that rows sharing a key stay mutually 
ordered. That invariant is preserved — the specific queue chosen is an internal 
detail, and every queue is backed by an equivalent writer for the table.
   
   Queue assignment is computed per row at write time and is not persisted in 
checkpoint state, so restoring a job taken before this change is unaffected.
   
   No config options, defaults, or SPI contracts changed, so no 
`incompatible-changes.md` entry is needed.
   
   ### How was this patch tested?
   
   Added 
`MultiTableSinkWriterTest.testMinValuePrimaryKeyHashRoutesToValidQueue`, which 
builds a writer with `multi_table_sink_replica = 3` and writes two rows — one 
with an `INT` primary key of `Integer.MIN_VALUE`, one with a `BIGINT` primary 
key of `Long.MIN_VALUE` — then asserts both rows reach a sink writer.
   
   I verified the test actually reproduces the bug by reverting the one-line 
fix and re-running it:
   
   ```
   [ERROR] testMinValuePrimaryKeyHashRoutesToValidQueue -- Time elapsed: 0.024 
s <<< FAILURE!
   org.opentest4j.AssertionFailedError: Unexpected exception thrown:
       java.lang.IndexOutOfBoundsException: Index -2 out of bounds for length 3
   ```
   
   With the fix restored, the full `seatunnel-api` suite passes: `Tests run: 
384, Failures: 0, Errors: 0, Skipped: 0`. `./mvnw spotless:apply` and `./mvnw 
-pl seatunnel-api -DskipTests verify` are also clean.
   
   ### Check list
   
   * [ ] If any new Jar binary package adding in your PR, please add License 
Notice according
     [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
   * [ ] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   * [ ] If necessary, please update `incompatible-changes.md` to describe the 
incompatibility caused by this PR.
   * [ ] If you are contributing the connector code, please check that the 
following files are updated:
     1. Update 
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
 and add new connector information in it
     2. Update the pom file of 
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
     3. Add ci label in 
[label-scope-conf](https://github.com/apache/seatunnel/blob/dev/.github/workflows/labeler/label-scope-conf.yml)
     4. Add e2e testcase in 
[seatunnel-e2e](https://github.com/apache/seatunnel/tree/dev/seatunnel-e2e/seatunnel-connector-v2-e2e/)
     5. Update connector 
[plugin_config](https://github.com/apache/seatunnel/blob/dev/config/plugin_config)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to