diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java b/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java index cd45fed2a3a..8cb630f0800 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java @@ -237,8 +237,16 @@ public static TopicPublishInfo topicRouteData2TopicPublishInfo(final String topi if (route.getOrderTopicConf() != null && route.getOrderTopicConf().length() > 0) { String[] brokers = route.getOrderTopicConf().split(";"); for (String broker : brokers) { - String[] item = broker.split(":"); - int nums = Integer.parseInt(item[1]); + String[] item = broker.split(":", 2); + if (item.length != 2) { + continue; + } + int nums; + try { + nums = Integer.parseInt(item[1]); + } catch (NumberFormatException e) { + continue; + } for (int i = 0; i < nums; i++) { MessageQueue mq = new MessageQueue(topic, item[0], i); info.getMessageQueueList().add(mq); diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/factory/MQClientInstanceTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/factory/MQClientInstanceTest.java index 82b9080438f..831d4a012d1 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/factory/MQClientInstanceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/factory/MQClientInstanceTest.java @@ -251,6 +251,17 @@ public void testTopicRouteData2TopicPublishInfoWithOrderTopicConf() { assertEquals(4, actual.getMessageQueueList().size()); } + @Test + public void testTopicRouteData2TopicPublishInfoSkipsMalformedOrderTopicConf() { + TopicRouteData topicRouteData = createTopicRouteData(); + topicRouteData.setOrderTopicConf("missing-count;invalid:not-a-number;127.0.0.1:2"); + + TopicPublishInfo actual = MQClientInstance.topicRouteData2TopicPublishInfo(topic, topicRouteData); + + assertTrue(actual.isOrderTopic()); + assertEquals(2, actual.getMessageQueueList().size()); + } + @Test public void testTopicRouteData2TopicPublishInfoWithTopicQueueMappingByBroker() { TopicRouteData topicRouteData = createTopicRouteData(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java index 8f08c1df0e5..cbb997e913c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java @@ -109,14 +109,22 @@ private static List buildWrite(TopicRouteWrapper topicR if (StringUtils.isNotBlank(topicRoute.getOrderTopicConf())) { String[] brokers = topicRoute.getOrderTopicConf().split(";"); for (String broker : brokers) { - String[] item = broker.split(":"); + String[] item = broker.split(":", 2); + if (item.length != 2) { + continue; + } String brokerName = item[0]; String brokerAddr = topicRoute.getMasterAddr(brokerName); if (brokerAddr == null) { continue; } - int nums = Integer.parseInt(item[1]); + int nums; + try { + nums = Integer.parseInt(item[1]); + } catch (NumberFormatException e) { + continue; + } for (int i = 0; i < nums; i++) { AddressableMessageQueue mq = new AddressableMessageQueue( new MessageQueue(topicRoute.getTopicName(), brokerName, i), diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelectorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelectorTest.java index e44ed28f4a6..57e37345410 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelectorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelectorTest.java @@ -81,4 +81,14 @@ public void testWriteMessageQueue() { messageQueueSelector.selectOne(false); assertEquals(queue, messageQueueSelector.selectOne(false)); } -} \ No newline at end of file + + @Test + public void testWriteMessageQueueSkipsMalformedOrderTopicConf() { + topicRouteData.setOrderTopicConf("missing-count;" + BROKER_NAME + ":not-a-number;" + BROKER_NAME + ":2"); + + MessageQueueSelector messageQueueSelector = + new MessageQueueSelector(new TopicRouteWrapper(topicRouteData, TOPIC), false); + + assertEquals(2, messageQueueSelector.getQueues().size()); + } +}