diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java index f3c68eab5c4..68660acbab9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java @@ -96,7 +96,13 @@ public boolean createTopicOnBroker(String topic, int wQueueNum, int rQueueNum, L Set curBrokerAddr = new HashSet<>(); if (curBrokerDataList != null) { for (BrokerData brokerData : curBrokerDataList) { - curBrokerAddr.add(brokerData.getBrokerAddrs().get(MixAll.MASTER_ID)); + if (brokerData == null || brokerData.getBrokerAddrs() == null) { + continue; + } + String addr = brokerData.getBrokerAddrs().get(MixAll.MASTER_ID); + if (addr != null) { + curBrokerAddr.add(addr); + } } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java index cdfc7f7fc23..ffdbbd29c91 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java @@ -17,6 +17,7 @@ package org.apache.rocketmq.proxy.service.admin; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.Set; @@ -87,6 +88,30 @@ public void testCreateTopic() throws Exception { assertEquals(8, topicConfigArgumentCaptor.getValue().getReadQueueNums()); } + @Test + public void testCreateTopicOnBrokerSkipsMalformedCurrentBrokerData() throws Exception { + BrokerData malformedBrokerData = new BrokerData(); + + ArgumentCaptor addrArgumentCaptor = ArgumentCaptor.forClass(String.class); + ArgumentCaptor topicConfigArgumentCaptor = ArgumentCaptor.forClass(TopicConfig.class); + doNothing().when(mqClientAPIExt) + .createTopic(addrArgumentCaptor.capture(), anyString(), topicConfigArgumentCaptor.capture(), anyLong()); + + assertTrue(defaultAdminService.createTopicOnBroker( + "createTopic", + 7, + 8, + Collections.singletonList(malformedBrokerData), + createTopicRouteData(1).getBrokerDatas(), + false, + 0 + )); + + assertEquals(1, addrArgumentCaptor.getAllValues().size()); + assertEquals("127.0.0.1:10911", addrArgumentCaptor.getAllValues().get(0)); + assertEquals("createTopic", topicConfigArgumentCaptor.getValue().getTopicName()); + } + private TopicRouteData createTopicRouteData(int brokerNum) { TopicRouteData topicRouteData = new TopicRouteData(); for (int i = 0; i < brokerNum; i++) { @@ -100,4 +125,4 @@ private TopicRouteData createTopicRouteData(int brokerNum) { } return topicRouteData; } -} \ No newline at end of file +}