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]