mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #2962] Implement DefaultMQAdminExt::examineTopicConfig function
This commit is contained in:
@@ -208,7 +208,7 @@ public class DefaultMQAdminExt extends ClientConfig implements MQAdminExt {
|
||||
}
|
||||
|
||||
@Override
|
||||
public TopicConfig examineTopicConfig(String addr, String topic) {
|
||||
public TopicConfig examineTopicConfig(String addr, String topic) throws RemotingException, InterruptedException, MQBrokerException {
|
||||
return defaultMQAdminExtImpl.examineTopicConfig(addr, topic);
|
||||
}
|
||||
|
||||
|
||||
@@ -222,8 +222,9 @@ public class DefaultMQAdminExtImpl implements MQAdminExt, MQAdminExtInner {
|
||||
}
|
||||
|
||||
@Override
|
||||
public TopicConfig examineTopicConfig(String addr, String topic) {
|
||||
return null;
|
||||
public TopicConfig examineTopicConfig(String addr, String topic) throws RemotingException, InterruptedException, MQBrokerException {
|
||||
TopicConfigSerializeWrapper topicConfigSerializeWrapper = this.mqClientInstance.getMQClientAPIImpl().getAllTopicConfig(addr,timeoutMillis);
|
||||
return topicConfigSerializeWrapper.getTopicConfigTable().get(topic);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -92,7 +92,7 @@ public interface MQAdminExt extends MQAdmin {
|
||||
|
||||
SubscriptionGroupConfig examineSubscriptionGroupConfig(final String addr, final String group);
|
||||
|
||||
TopicConfig examineTopicConfig(final String addr, final String topic);
|
||||
TopicConfig examineTopicConfig(final String addr, final String topic) throws RemotingException, InterruptedException, MQBrokerException;
|
||||
|
||||
TopicStatsTable examineTopicStats(
|
||||
final String topic) throws RemotingException, MQClientException, InterruptedException,
|
||||
|
||||
@@ -34,6 +34,7 @@ import org.apache.rocketmq.client.exception.MQClientException;
|
||||
import org.apache.rocketmq.client.impl.MQClientAPIImpl;
|
||||
import org.apache.rocketmq.client.impl.MQClientManager;
|
||||
import org.apache.rocketmq.client.impl.factory.MQClientInstance;
|
||||
import org.apache.rocketmq.common.TopicConfig;
|
||||
import org.apache.rocketmq.common.admin.ConsumeStats;
|
||||
import org.apache.rocketmq.common.admin.OffsetWrapper;
|
||||
import org.apache.rocketmq.common.admin.TopicOffset;
|
||||
@@ -54,6 +55,7 @@ import org.apache.rocketmq.common.protocol.body.ProcessQueueInfo;
|
||||
import org.apache.rocketmq.common.protocol.body.ProducerConnection;
|
||||
import org.apache.rocketmq.common.protocol.body.QueueTimeSpan;
|
||||
import org.apache.rocketmq.common.protocol.body.SubscriptionGroupWrapper;
|
||||
import org.apache.rocketmq.common.protocol.body.TopicConfigSerializeWrapper;
|
||||
import org.apache.rocketmq.common.protocol.body.TopicList;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
|
||||
@@ -225,6 +227,14 @@ public class DefaultMQAdminExtTest {
|
||||
consumerRunningInfo.setStatusTable(new TreeMap<String, ConsumeStatus>());
|
||||
consumerRunningInfo.setSubscriptionSet(new TreeSet<SubscriptionData>());
|
||||
when(mQClientAPIImpl.getConsumerRunningInfo(anyString(), anyString(), anyString(), anyBoolean(), anyLong())).thenReturn(consumerRunningInfo);
|
||||
|
||||
TopicConfigSerializeWrapper topicConfigSerializeWrapper = new TopicConfigSerializeWrapper();
|
||||
topicConfigSerializeWrapper.setTopicConfigTable(new ConcurrentHashMap<String, TopicConfig>() {
|
||||
{
|
||||
put("topic_test_examine_topicConfig", new TopicConfig("topic_test_examine_topicConfig"));
|
||||
}
|
||||
});
|
||||
when(mQClientAPIImpl.getAllTopicConfig(anyString(),anyLong())).thenReturn(topicConfigSerializeWrapper);
|
||||
}
|
||||
|
||||
@AfterClass
|
||||
@@ -406,4 +416,10 @@ public class DefaultMQAdminExtTest {
|
||||
assertThat(subscriptionGroupWrapper.getSubscriptionGroupTable().get("Consumer-group-one").getGroupName()).isEqualTo("Consumer-group-one");
|
||||
assertThat(subscriptionGroupWrapper.getSubscriptionGroupTable().get("Consumer-group-one").isConsumeBroadcastEnable()).isTrue();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testExamineTopicConfig() throws MQBrokerException, RemotingException, InterruptedException {
|
||||
TopicConfig topicConfig = defaultMQAdminExt.examineTopicConfig("127.0.0.1:10911", "topic_test_examine_topicConfig");
|
||||
assertThat(topicConfig.getTopicName().equals("topic_test_examine_topicConfig"));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user