SEPURI-SAI-KRISHNA commented on PR #12349: URL: https://github.com/apache/seatunnel/pull/12349#issuecomment-5707143613
You're right, and I reproduced it. The retry-topic interaction is a real defect in my change and I missed it entirely. One correction on attribution first. I counted the 50 failure markers in the JDK 11 job: | cause | tests | |---|---| | `broker[broker-a] not exist` at `generateTestData:467` | 22 | | `ROCKETMQ-11` from `offsetTopics` | 14 | | `ROCKETMQ-09` via the `%RETRY%` route | 14 | 36 of those are the pre-existing #12322 flakiness that #12323 fixes, which this branch doesn't carry. One test settles it anyway: `testSourceRocketMqTextToConsole` fails 7 of 7, inside the submitted job, at `run:180` -> `setPartitionStartOffset:380` -> `listConsumerGroupOffsets:457` on `CODE: 17 %RETRY%SeaTunnel-Consumer-Group`, and it never failed once in the 12 daily `Schedule Backend` runs I sampled for #12322. ## Code 17 can be told apart from a route gap I went through the 4.9.4 bytecode, which is both what the connector pins and what the E2E image runs. `DefaultMQAdminExtImpl.examineConsumeStats` resolves its route only from `%RETRY%<group>`: offsets 0-7 are `getRetryTopic` then `examineTopicRouteInfo`, with no try/catch and no fallback. The requested topic is only a parameter to the later `getConsumeStats` broker call, which throws `MQBrokerException` rather than `MQClientException`, and `getTopicRouteInfoFromNameServer` always throws on code 17 instead of returning null. So any `MQClientException` carrying code 17 out of `examineConsumeStats` is the retry topic and nothing else, without having to parse the topic name out of the message. `RouteInfoManager` decides the other half. `BROKER_CHANNEL_EXPIRED_TIME` is 120000L, and `scanNotActiveBroker` closes the channel then calls `onChannelDestroy`, which sweeps `topicQueueTable` keyed by brokerName. Route loss is broker scoped rather than topic scoped, so a genuine gap takes `%RETRY%<group>` and the data topic together. The only thing that removes one topic's route on its own is `RouteInfoManager.deleteTopic`, which is deliberate admin action. That is why Option A unguarded would put the bug back. It reads a real gap as a cold start and rewinds exactly as dev does, leaving #12332 open underneath a PR that claims to close it. Option B is out for the reason you gave me on #12332: `topicExist` swallows code 17 and returns false (`RocketMqAdminUtil.java:196-198`), so it carries the same ambiguity. ## What I'll implement Option A with one guard. On code 17, probe the requested topic's route using the admin client already open: - resolves, so the name server is healthy and the retry topic is genuinely absent: contribute nothing for that topic and carry on, preserving today's cold start - also fails, so the gap is broker wide: throw, as this PR intends Plus the WARN line you asked for on the first branch. Two things make that first branch safe rather than a guess. `ClientManageProcessor.heartBeat` creates `%RETRY%<group>` at offsets 132-163, before `registerConsumer` at 200, guarded only on the subscription group config that `SubscriptionGroupManager` auto-creates by default; a group without one can't pull at all, since `PullMessageProcessor` rejects it with `subscription group [%s] does not exist`. So committed offsets imply the retry topic was created. Separately, `setPartitionStartOffset` has a single call site at `run():180`, immediately after `fetchPendingPartitionSplit():179` under the same lock, and `discoverySplits():280` never calls it; since `offsetTopics` resolves the requested topic through `examineTopicStats` -> `examineTopicRouteInfo`, that route was live microseconds before the probe runs. Two cases it still won't cover, both behaving exactly as dev does today: someone deleting `%RETRY%<group>` out from under a live group, and a scale-out window where a new broker carries the data topic before the consumer has heartbeated to it. ## E2E and CI I said in the description that this couldn't be injected deterministically, which was too strong, since per-topic deletion does exist and dropping only the retry topic would do it. But `deleteTopicIfExist` (`RocketMqIT.java:903-925`) doesn't currently work: it passes broker names where addresses belong and the literal `"delete_topic"` where the cluster name belongs, and the blanket catch hides the result. In this job it logs `connect to null failed` 14 times and `Deleted topic` zero times. I'll raise that on its own rather than widen this PR, and cover both branches here with unit tests on the existing seam, which has no timing window. On wanting the connector E2E green here before merge: that can't happen while #12323 is unmerged, since the 36 failures above are the ones it fixes. The two PRs touch disjoint files, so I'll rebase onto that branch and note it in the description. If you'd still rather have Option A unguarded, say so and I'll cut the probe. Thanks for catching this. -- 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]
