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

   ### Purpose of this pull request
   
   Closes #12332.
   
   `RocketMqAdminUtil.currentOffsets` treated a failed group-offset lookup as 
proof that a consumer group had committed nothing. With `start.mode = 
CONSUME_FROM_GROUP_OFFSETS` that made the source silently rewind to the first 
offset and re-deliver the whole topic, while the committed offsets were intact 
and merely unreadable for a moment.
   
   The swallow was:
   
   ```java
   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, and 17 is also what the name server returns for a topic that 
exists but whose route is momentarily unavailable. The only production caller 
is `RocketMqSourceSplitEnumerator.listConsumerGroupOffsets`, feeding a branch 
where an empty map means "start from the beginning".
   
   **The contract this implements**, as agreed on the issue: an empty map is 
returned only when `examineConsumeStats` completed normally and the requested 
queues genuinely have no matching committed offsets. Any failed lookup, code 17 
included, now surfaces as `RocketMqConnectorException`.
   
   This does not break a real cold start, and the ordering is what guarantees 
it. In the enumerator, `RocketMqSourceSplitEnumerator.java:177-180`:
   
   ```java
   fetchPendingPartitionSplit();   // -> getTopicInfo() -> 
RocketMqAdminUtil.offsetTopics(...)
   setPartitionStartOffset();      // -> listConsumerGroupOffsets() -> 
currentOffsets(...)
   ```
   
   `offsetTopics` has no `TOPIC_NOT_EXIST` branch and raises 
`GET_MIN_AND_MAX_OFFSETS_ERROR` for a topic it cannot resolve, so a genuinely 
absent topic fails there first and never reaches the group-offset lookup. The 
branch this patch changes is therefore only reachable from a transient failure 
on that path.
   
   Changes:
   
   - `currentOffsets` no longer maps any exception onto an empty map. 
`InterruptedException` re-interrupts the current thread before being wrapped.
   - A package-private overload taking an already-started `DefaultMQAdminExt` 
is added as a test seam, so `examineConsumeStats` can be driven 
deterministically. No mutable static state, and the public signature is 
unchanged.
   - No retry policy and no new option, per the agreed scope.
   
   The legitimate empty cases are deliberately preserved: a successful lookup 
whose queues do not match the requested set, which is what happens for a newly 
discovered queue that has no consume-stats entry yet, still returns an empty 
map and still cold-starts that queue.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes, and it is worth stating plainly rather than filing this as 
internal-only.
   
   A job using `CONSUME_FROM_GROUP_OFFSETS` that previously survived a 
transient metadata failure by silently restarting from the first offset will 
now fail with `RocketMqConnectorException` carrying 
`GET_CONSUMER_GROUP_OFFSETS_ERROR` instead. That is the intended trade: a 
visible failure the operator can retry is better than an invisible full-topic 
replay into the sink. No configuration option, default value, or public API 
changes.
   
   ### How was this patch tested?
   
   `./mvnw -pl seatunnel-connectors-v2/connector-rocketmq test` on JDK 11: **9 
tests, 0 failures, 0 errors**, covering the three existing enumerator tests 
plus the new coverage.
   
   New `RocketMqAdminUtilTest`, driving `examineConsumeStats` through the seam:
   
   - a `TOPIC_NOT_EXIST` (code 17) failure surfaces as 
`RocketMqConnectorException` rather than an empty map
   - a successful lookup matching none of the requested queues still returns an 
empty map
   - a successful lookup that does match returns the committed offset unchanged
   
   New 
`RocketMqSourceSplitEnumeratorTest.testRun_doesNotFallBackToFirstOffsetWhenGroupOffsetLookupFails`:
 a failing group-offset lookup propagates, `flatOffsetTopics` (the first-offset 
fallback) is never called, and no split is assigned.
   
   Your three requested assertions map onto the suite as follows. The 
successful filtered empty case was already covered by the existing 
`testRun_fallsBackToFirstOffsetWhenGroupOffsetsAreMissing`, and the restore 
path by the existing `testRestore_preservesCheckpointOffsets` and 
`testRestore_newQueueUsesConfiguredStartMode`; all three still pass unchanged, 
which is the evidence that neither the acknowledgement nor the restore path 
moved. The failure case is the one genuinely new test.
   
   I also checked the new tests actually detect the bug rather than merely 
passing: with the old `TOPIC_NOT_EXIST` branch temporarily restored, 
`testCurrentOffsets_topicNotExistSurfacesAsFailure` fails and the other two 
still pass, so the coverage is pinned to the behaviour that changed.
   
   No E2E regression is included, and the reason is deliberate. The scenario 
can be injected, since `RocketMqIT.deleteTopicIfExist` already separates 
`deleteTopicInBroker` from `deleteTopicInNameServer` and calling only the 
latter leaves broker data and committed offsets intact while removing the 
route. But brokers re-register their topics with the name server on a periodic 
heartbeat, so the route returns on its own and such a test would be racing that 
interval. Given that I have spent #12323 removing exactly that class of timing 
dependence from this connector's e2e suite, adding one back here seemed like 
the wrong trade. Happy to add it with the caveat documented if you would prefer 
the coverage.
   
   One judgement call worth flagging: `RocketMqAdminUtil` had no test class, so 
unit-covering the new contract meant adding one, which sits against the "do not 
create new test classes casually" guidance. Putting `RocketMqAdminUtil` 
assertions inside the enumerator's test class seemed worse. Happy to move them.
   


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