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

   ### 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
   
   `RocketMqAdminUtil.currentOffsets` treats a transient metadata lookup 
failure as proof that a consumer group has committed nothing. With `start.mode 
= CONSUME_FROM_GROUP_OFFSETS`, that makes the source silently rewind to the 
first offset and re-deliver the whole topic, even though the committed offsets 
exist and are intact.
   
   **The swallow.** `RocketMqAdminUtil.java:298-307`:
   
   ```java
   } catch (MQClientException | MQBrokerException | RemotingException | 
InterruptedException e) {
       if (e instanceof MQClientException) {
           if (((MQClientException) e).getResponseCode() == 
ResponseCode.TOPIC_NOT_EXIST) {
               return Collections.emptyMap();
           }
   ```
   
   `ResponseCode.TOPIC_NOT_EXIST` is `17` in RocketMQ 4.9.4, the version this 
connector pins. That is the same code carried by `MQClientException: CODE: 17 
DESC: No topic route info in name server for the topic: <topic>`, which is 
raised when the name server has no route for a topic, not only when the topic 
is genuinely absent. So a route that is briefly unavailable and a topic that 
never existed are indistinguishable to this method, and both return an empty 
map.
   
   **Why that is load-bearing.** The only production caller is 
`RocketMqSourceSplitEnumerator.listConsumerGroupOffsets:455-458`, reached from 
`setPartitionStartOffset`, `RocketMqSourceSplitEnumerator.java:379-386`:
   
   ```java
   case CONSUME_FROM_GROUP_OFFSETS:
       Map<MessageQueue, Long> groupOffsets = listConsumerGroupOffsets(queues);
       if (groupOffsets.isEmpty()) {
           topicPartitionOffsets.putAll(listOffsets(queues, 
ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET));
       } else {
           topicPartitionOffsets.putAll(groupOffsets);
       }
   ```
   
   An empty map means "start from the beginning". The fallback itself is 
correct for a genuine cold start. The defect is that it is also reached when 
the answer was simply unreadable for a moment.
   
   **These route gaps are real and not brief.** While working on #12322 I 
measured one directly in CI, fork run `34972703578`, job `104395843312`. The 
name server returned no route for an existing, actively used topic 10 
consecutive times at roughly 30 second intervals, from 13:20:56 to 13:25:20, a 
continuous window of 4 minutes 24 seconds. Anything calling `currentOffsets` 
inside that window gets an empty map. `partition.discovery.interval.millis` 
makes the enumerator re-enter this path repeatedly, so a job only has to be 
unlucky once.
   
   ### Relationship to #9940 and #10778
   
   This is a third path to the symptom reported in #9940, and it is not covered 
by the two outcomes recorded there:
   
   - #10778 fixed the checkpoint restore path, where restored offsets were 
overwritten before the start mode was reapplied. This defect needs no restore 
and no failover; it happens on an ordinary start or on any subsequent split 
discovery.
   - The other outcome recorded in #9940 is that a cold start with no 
broker-side committed offset legitimately follows RocketMQ group-offset 
behavior. Agreed, and that is not this. Here the committed offsets exist on the 
broker and are simply not readable at that instant.
   
   Reported separately rather than on #9940 because the defect is a specific, 
separable piece of error handling with its own fix, and #9940 is a user report 
whose recorded conclusions both remain correct.
   
   ### What I have and have not verified
   
   Verified: the code paths above on current `dev`; that 
`ResponseCode.TOPIC_NOT_EXIST` is 17 in RocketMQ 4.9.4; that `CODE: 17` is what 
a missing name-server route produces, from CI logs; and that such gaps persist 
for minutes.
   
   Not verified: I have not captured an end-to-end reproduction of a production 
job rewinding to offset 0 through this specific path. The defect is established 
by construction rather than by a captured trace, and the #9940 symptom is 
consistent with it but is not proof of it. I would rather state that plainly 
than overclaim.
   
   ### Possible fix
   
   The two cases produce the same response code, so `currentOffsets` cannot 
distinguish them from the exception alone. `RocketMqAdminUtil` already has 
`topicExist(config, topic)` at line 186, which answers exactly the question 
that separates them. One option is to consult it before concluding empty: if 
the topic does exist, the lookup failed transiently and should be retried or 
raised rather than reported as "nothing committed"; if it genuinely does not 
exist, the current empty-map answer is right and the cold-start fallback is 
correct.
   
   Deliberately not proposing a patch yet, since choosing between retrying, 
failing fast, and distinguishing "unknown" from "known empty" in the return 
type is a semantics decision for maintainers. Happy to implement whichever 
direction is preferred.
   
   ### SeaTunnel Version
   
   `dev`, 3.0.0-SNAPSHOT. The same code is present in released versions 
carrying this connector.
   
   ### SeaTunnel Config
   
   ```conf
   source {
     Rocketmq {
       name.srv.addr = "..."
       topics = "..."
       start.mode = "CONSUME_FROM_GROUP_OFFSETS"
       consumer.group = "..."
       partition.discovery.interval.millis = "1000"
     }
   }
   ```
   
   ### Running Command
   
   ```shell
   sh bin/seatunnel.sh --config config/rocketmq_group_offsets.conf -e local
   ```
   
   ### Error Exception
   
   ```log
   org.apache.rocketmq.client.exception.MQClientException: CODE: 17  DESC: No 
topic route info in name server for the topic: <topic>
   ```
   
   No exception reaches the job. The failure is silent: the offsets are 
reported as empty and the source rewinds.
   
   ### Zeta or Flink or Spark Version
   
   Engine independent. The defect is in the connector's split enumerator and 
admin utility.
   
   ### Java or Scala Version
   
   Java 8 and Java 11.
   
   ### 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