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]