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]
