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

   ### Purpose of this pull request
   
   Fixes #12647.
   
   `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, discarding every offset already collected for the earlier topics 
in the list.
   
   `RocketMqAdminUtil.java:333-366` on `dev` at `11e4f5554`:
   
   ```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
           }
   ```
   
   The return rested on an invariant: a group cannot commit an offset without 
first registering, and registering creates the retry topic, so a missing retry 
topic should mean the accumulated map is still empty. The guard at `:339` 
probes **the requested topic**, not the group's retry topic. On a multi-broker 
cluster where the retry topic and a later topic live on different brokers, 
losing the retry topic's broker part way through the loop leaves the later 
topic resolvable, the guard passes, and offsets read successfully moments 
earlier are discarded.
   
   `RocketMqSourceSplitEnumerator` is the only production caller 
(`listConsumerGroupOffsets:455`), and its emptiness check is all or nothing 
(`:379-384`), falling back to `listOffsets(..., CONSUME_FROM_FIRST_OFFSET)` 
which assigns `getMinOffset()` to **every** queue (`:430`). So with `start.mode 
= CONSUME_FROM_GROUP_OFFSETS` and two or more topics, one topic's missing retry 
route replayed all of them.
   
   This change turns that `return` into a `continue`, which is the fix named in 
the source comment added by #12438. The unused `java.util.Collections` import 
goes with it, and the comment and `log.warn` text are updated to describe the 
per-topic contract rather than a whole-lookup one.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes, a behavioural fix on a failure path. No config option, default or 
public signature changes, so there is nothing to add to 
`incompatible-changes.md`.
   
   Nothing is worse off than today. A fresh split is constructed with 
`startOffset = topicOffset.getMinOffset()` (`getTopicInfo:313`), and 
`setPartitionStartOffset` only overwrites queues present in the map 
(`:412-417`):
   
   - before, empty map: every queue of every topic rewound to its minimum 
offset.
   - now: the earlier topics resume from their committed offsets, and the 
skipped topic's queues keep the minimum offset they were constructed with, 
which is exactly the position they get today.
   
   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, and 
`testCurrentOffsets_retryTopicMissingWhileRouteHealthyReturnsEmpty` pins that.
   
   Out of scope, and noted in the issue: skipping the topic does not make that 
topic's own replay correct in the route-loss case, it only stops the replay 
from spreading to the rest of the list. Separating "never registered" from "the 
retry topic's broker is gone" needs more than a probe of the requested topic, 
and that probe is a deliberate heuristic for cluster health.
   
   ### How was this patch tested?
   
   A new unit test beside 
`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.
   
   It fails on unmodified `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: <{}>
   ```
   
   and passes with the fix:
   
   ```
   ./mvnw -pl seatunnel-connectors-v2/connector-rocketmq test
   Tests run: 25, Failures: 0, Errors: 0, Skipped: 0
   ```
   
   Mutation checked: restoring `return Collections.emptyMap()` fails only this 
test, 1 of 7 in the class, so the new test is the sole guard on the behaviour.
   
   `spotless:check` is green on the module with no reformatting needed.
   
   ### Why no IT change
   
   `RocketMqIT` exists for this connector, but the trigger is a route asymmetry 
between the group's retry topic and a later topic in the list, which requires 
those two to be hosted on different brokers. The IT runs a single broker 
container, so the condition cannot be constructed there. All three 
`currentOffsets` call sites in `RocketMqIT` are single topic and are 
unaffected. Driving `examineConsumeStats` through a mocked `DefaultMQAdminExt`, 
which is what the package-private overload exists for, is the only level at 
which this is deterministic.
   


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