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]

Reply via email to