SEPURI-SAI-KRISHNA opened a new issue, #12322:
URL: https://github.com/apache/seatunnel/issues/12322

   ### Search before asking
   
   - [X] I had searched in the 
[issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22)
 and found no similar issues.
   
   ### What happened
   
   `RocketMqIT` is the single largest source of flaky CI failures on `dev`, and 
the failures are not ambient infrastructure noise. They come from one 
reproducible gap in the test fixture.
   
   I sampled the `Schedule Backend` workflow, the daily cron that runs against 
`dev` with no pull request diff involved, so any failure there is 
dev-intrinsic. Twelve runs, 2026-08-28 through 2026-09-14: 1126 job-runs, 62 
failed jobs, every failed job log parsed.
   
   `rocketmq-connector-it` fails on **8 of the 12 days**, a 6/12 failure rate 
on the JDK 8 leg and 4/12 on JDK 11. On 2026-09-14 the run was `Tests run: 123, 
Failures: 20, Errors: 41`, elapsed 3,827s. The failing methods fail on **all 
ten** `TestContainer` templates, so this is not a single-container flake.
   
   Ranked by distinct days failed out of 12, eight of the ten flakiest tests in 
the whole repository are in this one class:
   
   | Days | Test |
   |---|---|
   | 7 | `testSourceRocketMqStartConfig` |
   | 6 | `testSourceRocketMqTextTagToConsole` |
   | 6 | `testSourceRocketMqTextErrorTagToConsole` |
   | 6 | `testSourceRocketMqMultiTableToAssert` |
   | 6 | `testSinkRocketMq` |
   | 6 | `testRocketMqSpecificOffsetsToConsole` |
   | 6 | `testRocketMqEarliestToConsole` |
   | 2 | `testSourceRocketMqRestore` |
   
   ### Root cause
   
   Every failure resolves to the same runtime error, reached by two different 
paths:
   
   ```
   MQClientException: CODE: 17  DESC: No topic route info in name server
                                for the topic: test_topic_source
   ```
   
   **Path A, the producer side.** `generateTestData` writes with
   
   ```java
   producer.send(message, new MessageQueue(topic, 
RocketMqContainer.BROKER_NAME, 0));
   ```
   
   at `RocketMqIT.java:467`. Addressing a `MessageQueue` directly bypasses the 
route resolution that a plain `send(Message)` performs, so the send fails 
outright if the name server has not published a route for that topic. This 
accounts for 31 of the 61 bad results on 2026-09-14, surfacing as 
`...->generateTestData:467 » MQClientException` under 
`testSourceRocketMqTextTagToConsole`, 
`testSourceRocketMqTextErrorTagToConsole`, 
`testSourceRocketMqMultiTableToAssert` and `testSourceRocketMqStartConfig`.
   
   **Path B, the consumer side.** The submitted job hits the same missing route 
through the split enumerator:
   
   ```
   RocketMqConnectorException: ErrorCode:[ROCKETMQ-11],
     ErrorDescription:[Failed to get topic min and max topic]
       at RocketMqAdminUtil.offsetTopics(RocketMqAdminUtil.java:231)
       at 
RocketMqSourceSplitEnumerator.getTopicInfo(RocketMqSourceSplitEnumerator.java:304)
       at RocketMqSourceSplitEnumerator.fetchPendingPartitionSplit(...:290)
   Caused by: MQClientException: CODE: 17 DESC: No topic route info in name 
server
   ```
   
   This is how `testRocketMqEarliestToConsole` and 
`testRocketMqSpecificOffsetsToConsole` fail. Note that 
`fetchPendingPartitionSplit` calls `getTopicInfo` for **every** start mode, so 
`testRocketMqLatestToConsole` and `testRocketMqTimestampToConsole` are equally 
exposed and simply have not been executing inside the window where the route 
was absent.
   
   **The gap.** The file already contains the remedy, 
`waitForTopicRoute(topic)` at line 470, and `startUp()` carries a comment 
explaining precisely why it is needed. But it is only ever called at three 
points:
   
   - `startUp()`, for `test_topic_source`
   - line 198, for `test_topic`
   - line 217, for `test_text_topic`
   
   whereas `generateTestData` is called at lines 239, 257, 273, 293, 311, 333, 
353, 354 and 406 and warms nothing itself, and the four source tests at lines 
362, 370, 378 and 386 submit jobs without warming anything. The guard was 
written but never applied at the point of use.
   
   ### Secondary symptom, same origin
   
   ```
   testSinkRocketMq:205 -> getRocketMqConsumerData:558 -> 
waitConsumedOffsetsSynced:577
     » ConditionTimeout
   ```
   
   10 of the 61. The in-file comment records that this window was already 
raised from 30s to 60s; it still timed out at 60s on 2026-09-14, which suggests 
the ceiling is not the real constraint here either.
   
   ### Onset
   
   The failure shape changes sharply in the sample. Before 2026-09-04, across 8 
`rocketmq-connector-it` job-runs, there were 2 failures and each was a single, 
different method; 09-01, 09-02 and 09-03 were fully green on both JDK legs. 
From 2026-09-04 onward, across 14 job-runs, there were 7 failures and every one 
of them was the same seven-method cluster.
   
   That window contains #12050, which raised 
`RocketMqContainer.DEFAULT_TOPIC_QUEUE_NUMS` from 1 to 4 and feeds it to the 
broker as `defaultTopicQueueNums` with `autoCreateTopicEnable=true`. I want to 
be precise about the strength of this: the correlation in onset and shape is 
clear, but I have not proven the causal chain from queue count to route 
availability, so I am reporting it as onset evidence only. The fix proposed 
below does not depend on it, since it is justified by the missing-route error 
on its own.
   
   ### SeaTunnel Version
   
   `dev`, 3.0.0-SNAPSHOT
   
   ### SeaTunnel Config
   
   
`seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq/*.conf`,
 unchanged by this report.
   
   ### Running Command
   
   ```shell
   ./mvnw -B -pl 
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e -am verify \
     -DskipUT=true -DskipIT=false
   ```
   
   ### Error Exception
   
   ```log
   [ERROR] RocketMqIT.testSourceRocketMqStartConfig:406->generateTestData:467 » 
MQClient
   [ERROR] 
RocketMqIT.testSourceRocketMqTextTagToConsole:239->generateTestData:467 » 
MQClient
   [ERROR] 
RocketMqIT.testSourceRocketMqTextErrorTagToConsole:257->generateTestData:467 » 
MQClient
   [ERROR] 
RocketMqIT.testSourceRocketMqMultiTableToAssert:333->generateTestData:467 » 
MQClient
   [ERROR] 
RocketMqIT.testSinkRocketMq:205->getRocketMqConsumerData:558->waitConsumedOffsetsSynced:577
 » ConditionTimeout
   [ERROR] Tests run: 123, Failures: 20, Errors: 41, Skipped: 0
   ```
   
   ### Zeta or Flink or Spark Version
   
   All ten `TestContainer` templates, including Flink 1.15.3 / 1.18.0, Spark 
2.4.6 / 3.3.0 and Zeta.
   
   ### Java or Scala Version
   
   JDK 8 and JDK 11, ubuntu-latest. The JDK 8 leg fails more often (6/12 
against 4/12).
   
   ### Screenshots
   
   _No response_
   
   ### Are you willing to submit PR?
   
   - [X] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [X] I agree to follow this project's [Code of 
Conduct](https://www.apache.org/foundation/policies/conduct)
   


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