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]