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

   ### Purpose of this pull request
   
   Closes #12322.
   
   `RocketMqIT` is the largest single source of flaky CI failures on `dev`. 
Measured over the `Schedule Backend` daily cron, which runs against `dev` with 
no pull request diff involved, `rocketmq-connector-it` failed on 8 of 12 
sampled days between 2026-08-28 and 2026-09-14, and eight of the ten flakiest 
tests in the repository belong to this one class.
   
   Every failure resolves to the same error, `MQClientException: CODE: 17 DESC: 
No topic route info in name server`, reached two ways:
   
   - `generateTestData` writes with `producer.send(message, new 
MessageQueue(topic, BROKER_NAME, 0))` at line 467. Addressing a queue directly 
bypasses the route resolution a plain `send(Message)` performs, so the send 
fails if the name server has no route for that topic.
   - The submitted source job hits the same missing route through 
`RocketMqSourceSplitEnumerator.getTopicInfo` and 
`RocketMqAdminUtil.offsetTopics`, surfacing as `ROCKETMQ-11 Failed to get topic 
min and max topic`.
   
   The file already contains the remedy, `waitForTopicRoute(topic)`, and 
`startUp()` carried a comment explaining exactly why it is needed. It was only 
applied at three points: `startUp()` for `test_topic_source`, and the two sink 
tests for `test_topic` and `test_text_topic`. Meanwhile `generateTestData` is 
called from nine sites and warmed nothing itself, and the four source tests 
submitted jobs without warming anything. This patch applies the existing guard 
at the point of use instead.
   
   Changes:
   
   - `generateTestData` warms the route for the topic it is about to write, 
which covers all nine call sites. The rationale comment moves from `startUp()` 
to this method, where the guard now lives.
   - The four source tests warm `test_topic_source` before `executeJob`. All 
four are included deliberately: `fetchPendingPartitionSplit` calls 
`getTopicInfo` for every start mode, so `testRocketMqLatestToConsole` and 
`testRocketMqTimestampToConsole` are exposed on the same path as the two that 
have been failing, and have simply not been executing inside the window where 
the route was absent.
   - `startUp()` keeps its post-write route confirmation and no longer needs 
the pre-write one, since `generateTestData` now performs it.
   
   No test is skipped, no timeout is raised, and no assertion is weakened.
   
   One point of precision on scope. The failure shape changes sharply at 
2026-09-04: before it, 2 failures across 8 job-runs, each a single different 
method, with 09-01 through 09-03 fully green on both JDK legs; after it, 7 
failures across 14 job-runs, every one the same seven-method cluster. That 
window contains #12050, which raised `DEFAULT_TOPIC_QUEUE_NUMS` from 1 to 4. I 
have not proven a causal chain from queue count to route availability, so I am 
not claiming one and this patch does not touch that constant. The change here 
is justified by the missing-route error on its own.
   
   ### Does this PR introduce _any_ user-facing change?
   
   No. The change is confined to one e2e test class. No production code, 
connector behavior, configuration or documentation is affected.
   
   ### How was this patch tested?
   
   The diagnosis is built from CI evidence rather than a local reproduction, 
because the failure depends on name server route timing inside the e2e 
container and does not reproduce on demand.
   
   - Sampled 12 `Schedule Backend` runs, 2026-08-28 to 2026-09-14: 1126 
job-runs, 62 failed jobs, every failed job log parsed. Every job that failed 
also passed on another day, so none of this is a permanently broken job.
   - Confirmed the root error from the full 2026-09-14 `rocketmq-connector-it 
(8, ubuntu-latest)` log on both paths described above, including the enumerator 
stack trace.
   - Verified against `RocketMqSourceSplitEnumerator` that `getTopicInfo` is 
reached for every start mode, which is why all four source tests are warmed 
rather than only the two observed failing.
   - Verified that `RocketMqSourceSplitEnumerator` skips queues with no 
specified offset rather than failing, which rules out the single-queue test 
configs as a contributing cause. The `.conf` files are deliberately untouched.
   - `./mvnw spotless:apply` runs clean on the module with JDK 11. Since 
google-java-format has to
     parse the file to reformat it, that also confirms the edit is 
syntactically valid.
   - A full local reactor build did not complete in my environment, and I want 
to be straight about
     why rather than imply more verification than I did: it fails earlier in 
`seatunnel-common`
     test sources against `org.apache.seatunnel.shade.com.typesafe.config`, 
because the shaded
     SNAPSHOT artifacts are not installed in my local repository. That is 
unrelated to this patch
     and the rocketmq module is never reached. The change itself only adds 
calls to
     `waitForTopicRoute(String)`, an existing private method in the same class, 
and removes one
     now-redundant call plus a comment, so it introduces no new symbol or 
signature. Compilation is
     covered by this PR's own CI.
   
   The real verification is the `rocketmq-connector-it` legs of this PR's own 
CI run, plus the `Schedule Backend` runs after it merges. Given a 6/12 base 
rate, a few consecutive green days is the signal to watch for.
   


-- 
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