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]