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

   ### 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` loops over the configured topic list and 
merges each topic's consume stats. When one topic answers `TOPIC_NOT_EXIST` 
while that topic's own route still resolves, the method returns an empty map 
immediately, throwing away every offset it has already collected for the 
earlier topics in the list. The only production caller reads an empty map as 
"this group has committed nothing" and rewinds **every** topic to its first 
offset, so a route problem affecting one topic re-delivers all of them.
   
   **What I verified and what I did not.**
   
   Verified, on `dev` at `11e4f5554`: the code paths quoted below; that the 
discard is reachable; that the caller rewinds on an empty map; and the 
behaviour change, by a unit test that fails on current `dev` and passes with 
the one-line fix. The failure output is included.
   
   Not verified: I have not captured this happening on a live multi-broker 
cluster. The trigger is a route asymmetry between the group's retry topic and a 
later topic in the list, which I can describe exactly but cannot schedule on 
demand. The defect is established from the code and from a unit test, not from 
a production trace.
   
   ### The discard
   
   `RocketMqAdminUtil.java:333-366`:
   
   ```java
   for (String topic : topics) {
       try {
           ConsumeStats consumeStats = adminClient.examineConsumeStats(groupId, 
topic);
           consumerOffsets.putAll(consumeStats.getOffsetTable());
       } catch (MQClientException e) {
           if (e.getResponseCode() == ResponseCode.TOPIC_NOT_EXIST
                   && topicRouteAvailable(adminClient, topic)) {
               ...
               return Collections.emptyMap();   // line 366
           }
   ```
   
   `consumerOffsets` is accumulated across iterations, so returning on 
iteration *n* drops iterations *1..n-1*. `topics` is a supported multi-topic 
list (`RocketMqSourceOptions.TOPICS` accepts a comma separated list).
   
   ### Why the return is not safe
   
   The return rests on an invariant: a group cannot commit an offset for any 
topic without first registering, and registering is what creates the retry 
topic, so a missing retry topic should mean nothing was committed for any topic 
and the accumulated map should still be empty.
   
   That invariant holds for the cold start it was written for, but the guard at 
line 339 probes **the requested topic**, not the group's retry topic, and the 
retry topic is re-resolved on every iteration. On a multi-broker cluster where 
the retry topic and a later topic in the list live on different brokers, losing 
the retry topic's broker part way through the loop leaves the later topic 
resolvable. The guard passes, and the return discards offsets that were read 
successfully moments earlier.
   
   This was recorded in the source comment at that branch by #12438, which also 
named the fix. The early return itself was first documented by #12417.
   
   ### Why it matters
   
   `RocketMqSourceSplitEnumerator` is the only production caller 
(`listConsumerGroupOffsets:455`). `RocketMqSourceSplitEnumerator.java:379-384`:
   
   ```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);
       }
   ```
   
   The emptiness check is all or nothing, and `listOffsets(..., 
CONSUME_FROM_FIRST_OFFSET)` assigns `topicOffset.getMinOffset()` to **every** 
queue (`:430`). So with `start.mode = CONSUME_FROM_GROUP_OFFSETS` and two or 
more topics, one topic's missing retry route replays all of them.
   
   ### Blast radius, before and after
   
   A fresh split is built with `startOffset = topicOffset.getMinOffset()` 
(`getTopicInfo:313`), and `setPartitionStartOffset` only overwrites queues 
present in the map (`:412-417`). So keeping a partial map does not leave the 
skipped topic undefined:
   
   - today, empty map: every queue of every topic rewinds to its minimum offset.
   - with the offsets kept: the earlier topics resume from their committed 
offsets, and the skipped topic's queues keep the minimum offset they were 
constructed with, which is the same answer it gets today.
   
   No queue is worse off, and the topics whose offsets were read correctly stop 
being replayed.
   
   ### Known residual
   
   Skipping the topic does not make its own replay correct in the route-loss 
case; it only stops that replay from spreading to the rest of the list. 
Distinguishing "this group never registered" from "the retry topic's broker is 
gone" needs more than a probe of the requested topic, and the requested-topic 
probe is a deliberate heuristic for cluster health. I am treating that as out 
of scope here rather than changing the heuristic.
   
   ### Reproduction
   
   A unit test beside 
`RocketMqAdminUtilTest.testCurrentOffsets_multipleTopicsAreMergedAcrossTheLoop`:
 the first topic returns consume stats, the second raises `TOPIC_NOT_EXIST` 
with `examineTopicRouteInfo` for that second topic returning a route, so the 
guard passes. On current `dev`:
   
   ```
   [ERROR] Tests run: 7, Failures: 1, Errors: 0, Skipped: 0
   [ERROR] 
RocketMqAdminUtilTest.testCurrentOffsets_laterTopicFailureKeepsOffsetsFromEarlierTopics:260
   the offset already collected for the first topic must survive a later 
topic's missing retry route
   ==> expected: <{MessageQueue [topic=test-topic, brokerName=broker-a, 
queueId=0]=42}> but was: <{}>
   ```
   
   With `return Collections.emptyMap()` changed to `continue`, all 25 tests in 
`connector-rocketmq` pass, and reverting the change fails only this one test.
   
   The cold start contract is unchanged: when every topic answers that way the 
map still ends up empty, which is the answer the caller needs. 
`testCurrentOffsets_retryTopicMissingWhileRouteHealthyReturnsEmpty` pins that 
and still passes.
   
   ### Not a duplicate
   
   - #12332 (closed) and its fix #12349 made a lost route surface as a failure 
instead of as "nothing committed". That work introduced the guard this report 
is about. It did not address the accumulated offsets being discarded.
   - #12383 is a broker-side offset visibility problem in `RocketMqIT` on a 
single topic (`test_topic`). This needs two or more topics and a route 
asymmetry, so the two do not overlap, though both touch `currentOffsets`.
   - #9940 and its fix #10778 are the checkpoint restore path, where restored 
offsets were overwritten before the start mode was reapplied. This needs no 
restore and no failover.
   
   The other two `TOPIC_NOT_EXIST` sites in the same class are correct and are 
not part of this: `topicExist:189` handles a single topic with no accumulation, 
and `offsetTopics:219` loops but throws on failure rather than discarding.
   
   ### SeaTunnel Version
   
   `dev` at `11e4f5554`, 3.0.0-SNAPSHOT.
   
   The offsets are also lost at the `v3.0.0` tag, by a different shape: there 
the `try` wraps the whole loop (`RocketMqAdminUtil.java:286-289`), so any 
`TOPIC_NOT_EXIST` aborts the loop and the catch at `:303-304` returns an empty 
map with no route guard at all. The per-iteration catch and the route probe 
arrived later with #12349. So this is not a regression introduced by that work; 
the discard survived the restructuring.
   
   ### SeaTunnel Config
   
   ```conf
   source {
     Rocketmq {
       name.srv.addr = "..."
       topics = "topic_a,topic_b"
       start.mode = "CONSUME_FROM_GROUP_OFFSETS"
       consumer.group = "my_group"
       format = "json"
     }
   }
   ```
   
   ### Are you willing to submit PR?
   
   - [X] Yes I am willing to submit a PR!
   


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