SEPURI-SAI-KRISHNA opened a new pull request, #12393:
URL: https://github.com/apache/seatunnel/pull/12393

   ### Purpose of this pull request
   
   `RocketMqIT.deleteTopicIfExist` has never deleted a topic. I found this 
while working on #12349 and said there that I would raise it separately.
   
   Three defects compound:
   
   * `deleteTopicInBroker(Set<String> addrs, String topic)` and 
`deleteTopicInNameServer(Set<String> addrs, String topic, String clusterName)` 
both expect **addresses**. The helper passed broker **names**, collected from 
`QueueData::getBrokerName`. A broker name does not resolve to a host, so the 
request fails immediately.
   * The literal `"delete_topic"` was passed as the third argument, which is 
the cluster name rather than a reason string. Confirmed against 4.9.4: in 
`MQClientAPIImpl.deleteTopicInNameServer` the second argument reaches 
`setTopic` and the third reaches `setClusterName`.
   * A blanket `catch (Exception e) { log.warn(...) }` reduced the whole thing 
to a warning nobody reads in a passing run.
   
   The first defect fails first, so the name server call is never even reached.
   
   **Evidence:** `Failed to delete topic ...: connect to null failed` appears 
**14 times in each of three job logs**, across two separate CI runs of this 
suite, and `Deleted topic` appears **zero** times in any of them. The callers, 
`testSourceRocketMqTextTagToConsole` and 
`testSourceRocketMqTextErrorTagToConsole`, therefore start from whatever state 
the previous run left behind rather than from the clean topic they assume.
   
   ### The fix
   
   Broker addresses now come from `TopicRouteData.getBrokerDatas()`, taking 
every entry of `getBrokerAddrs()` rather than one per broker. That matches 
RocketMQ's own `DeleteTopicSubCommand`, which resolves masters and slaves and 
deletes against all of them. `selectBrokerAddr()` would have been the obvious 
single-value choice, but it returns the master only when one is registered and 
otherwise picks at random, which is not a sound basis for a delete. The name 
server call passes a null address set, which makes the client resolve the name 
server list itself, and a null cluster name, which makes the name server take 
its global `RouteInfoManager.deleteTopic(topic)` branch instead of scoping the 
request to a cluster that never existed. I checked that branch rather than 
assuming it: `DefaultRequestProcessor.deleteTopicInNamesrv` tests the cluster 
name for null and for empty, and falls through to the global delete in both 
cases.
   
   A failure now throws instead of being logged and dropped. A helper named 
`deleteTopicIfExist` that returns normally should mean the topic is gone; if it 
cannot guarantee that, the caller is asserting against state it did not 
establish. The "topic is not there" case is still handled explicitly and 
returns quietly, so this only throws on a real failure.
   
   I also dropped a dead `admin != null` check in the `finally`, since `admin` 
is assigned from `new DefaultMQAdminExt()` on the line above and cannot be null.
   
   ### Two things I want to be straight about
   
   **This is the first time the deletion will actually happen.** Until now 
these two tests have been running against an undeleted topic. Making the 
cleanup work is a genuine behaviour change for them, and this PR's own CI is 
the first real signal on whether they are happy with a topic that is truly 
reset. I would rather find that out here than leave a cleanup helper that lies.
   
   **I am not claiming this fixes their flakiness.** Both callers sit in the 
`RocketMqIT` failure cluster from #12322 at 6 of 12 sampled days, and stale 
topic state is a plausible contributor, but I have not demonstrated that 
connection and I am not asserting it.
   
   ### Does this PR introduce _any_ user-facing change?
   
   No. Test-only, one E2E class, no production code.
   
   ### How was this patch tested?
   
   `./mvnw -q -DskipTests verify -pl 
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e` on JDK 11 
passes, covering the enforcer checks, `spotless:check` and compilation. 
`spotless:check` was also run separately with up-to-date caching disabled.
   
   No new test class is added: the change is to `RocketMqIT` itself, so the IT 
is both the subject and the coverage, and the `rocketmq-connector-it` legs 
exercise it through the two callers above.
   
   ### One deferred item
   
   @DanielLeens asked on #12349 for a comment on the multi-topic early-return 
in `RocketMqAdminUtil.currentOffsets`, and agreed it could ride along with this 
change. I have left it out deliberately: it sits in a method #12349 is 
currently rewriting, so adding it here would put that PR into conflict. It will 
follow once #12349 merges.
   


-- 
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