mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
This commit is contained in:
@@ -138,6 +138,7 @@ import org.apache.rocketmq.common.protocol.header.filtersrv.RegisterMessageFilte
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.AddWritePermOfBrokerRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.AddWritePermOfBrokerResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.DeleteKVConfigRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.DeleteTopicFromNamesrvRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.GetKVConfigRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.GetKVConfigResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.GetKVListByNamespaceRequestHeader;
|
||||
@@ -1446,7 +1447,7 @@ public class MQClientAPIImpl {
|
||||
}
|
||||
|
||||
public void deleteTopicInBroker(final String addr, final String topic, final long timeoutMillis)
|
||||
throws RemotingException, MQBrokerException, InterruptedException, MQClientException {
|
||||
throws RemotingException, InterruptedException, MQClientException {
|
||||
DeleteTopicRequestHeader requestHeader = new DeleteTopicRequestHeader();
|
||||
requestHeader.setTopic(topic);
|
||||
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.DELETE_TOPIC_IN_BROKER, requestHeader);
|
||||
@@ -1466,8 +1467,8 @@ public class MQClientAPIImpl {
|
||||
}
|
||||
|
||||
public void deleteTopicInNameServer(final String addr, final String topic, final long timeoutMillis)
|
||||
throws RemotingException, MQBrokerException, InterruptedException, MQClientException {
|
||||
DeleteTopicRequestHeader requestHeader = new DeleteTopicRequestHeader();
|
||||
throws RemotingException, InterruptedException, MQClientException {
|
||||
DeleteTopicFromNamesrvRequestHeader requestHeader = new DeleteTopicFromNamesrvRequestHeader();
|
||||
requestHeader.setTopic(topic);
|
||||
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.DELETE_TOPIC_IN_NAMESRV, requestHeader);
|
||||
|
||||
@@ -1485,7 +1486,7 @@ public class MQClientAPIImpl {
|
||||
}
|
||||
|
||||
public void deleteSubscriptionGroup(final String addr, final String groupName, final boolean removeOffset, final long timeoutMillis)
|
||||
throws RemotingException, MQBrokerException, InterruptedException, MQClientException {
|
||||
throws RemotingException, InterruptedException, MQClientException {
|
||||
DeleteSubscriptionGroupRequestHeader requestHeader = new DeleteSubscriptionGroupRequestHeader();
|
||||
requestHeader.setGroupName(groupName);
|
||||
requestHeader.setRemoveOffset(removeOffset);
|
||||
|
||||
+1
-1
@@ -20,7 +20,7 @@ import org.apache.rocketmq.remoting.CommandCustomHeader;
|
||||
import org.apache.rocketmq.remoting.annotation.CFNotNull;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
|
||||
public class DeleteTopicInNamesrvRequestHeader implements CommandCustomHeader {
|
||||
public class DeleteTopicFromNamesrvRequestHeader implements CommandCustomHeader {
|
||||
@CFNotNull
|
||||
private String topic;
|
||||
|
||||
+3
-3
@@ -39,7 +39,7 @@ import org.apache.rocketmq.common.protocol.body.RegisterBrokerBody;
|
||||
import org.apache.rocketmq.common.protocol.body.TopicConfigSerializeWrapper;
|
||||
import org.apache.rocketmq.common.protocol.header.GetTopicsByClusterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.DeleteKVConfigRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.DeleteTopicInNamesrvRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.DeleteTopicFromNamesrvRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.GetKVConfigRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.GetKVConfigResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.GetKVListByNamespaceRequestHeader;
|
||||
@@ -438,8 +438,8 @@ public class DefaultRequestProcessor extends AsyncNettyRequestProcessor implemen
|
||||
private RemotingCommand deleteTopicInNamesrv(ChannelHandlerContext ctx,
|
||||
RemotingCommand request) throws RemotingCommandException {
|
||||
final RemotingCommand response = RemotingCommand.createResponseCommand(null);
|
||||
final DeleteTopicInNamesrvRequestHeader requestHeader =
|
||||
(DeleteTopicInNamesrvRequestHeader) request.decodeCommandCustomHeader(DeleteTopicInNamesrvRequestHeader.class);
|
||||
final DeleteTopicFromNamesrvRequestHeader requestHeader =
|
||||
(DeleteTopicFromNamesrvRequestHeader) request.decodeCommandCustomHeader(DeleteTopicFromNamesrvRequestHeader.class);
|
||||
|
||||
this.namesrvController.getRouteInfoManager().deleteTopic(requestHeader.getTopic());
|
||||
|
||||
|
||||
Reference in New Issue
Block a user