diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java index 75f7089c5e0..16e5524fb4e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java @@ -78,7 +78,7 @@ public CompletableFuture queryRoute(ProxyContext ctx, QueryR String brokerName = queueData.getBrokerName(); Map brokerIdMap = brokerMap.get(brokerName); if (brokerIdMap == null) { - break; + continue; } for (Broker broker : brokerIdMap.values()) { messageQueueList.addAll(this.genMessageQueueFromQueueData(queueData, request.getTopic(), topicMessageType, broker)); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java index abbf82452ef..85c1b17ecb2 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java @@ -119,6 +119,29 @@ public void testQueryRoute() throws Throwable { } } + @Test + public void testQueryRouteSkipsMissingBrokerAndContinues() throws Throwable { + when(this.messagingProcessor.getTopicRouteDataForProxy(any(), any(), anyString())) + .thenReturn(createProxyTopicRouteDataWithMissingBrokerBeforeValidBroker()); + when(this.messagingProcessor.getMetadataService().getTopicMessageType(any(), anyString())) + .thenReturn(TopicMessageType.NORMAL); + + QueryRouteResponse response = this.routeActivity.queryRoute( + createContext(), + QueryRouteRequest.newBuilder() + .setEndpoints(grpcEndpoints) + .setTopic(Resource.newBuilder().setName(TOPIC).build()) + .build() + ).get(); + + assertEquals(Code.OK, response.getStatus().getCode()); + assertEquals(4, response.getMessageQueuesCount()); + for (MessageQueue messageQueue : response.getMessageQueuesList()) { + assertEquals(BROKER_NAME, messageQueue.getBroker().getName()); + assertEquals(grpcEndpoints, messageQueue.getBroker().getEndpoints()); + } + } + @Test public void testQueryRouteTopicExist() throws Throwable { when(this.messagingProcessor.getTopicRouteDataForProxy(any(), any(), anyString())) @@ -226,6 +249,14 @@ private static ProxyTopicRouteData createProxyTopicRouteData(int r, int w, int p return proxyTopicRouteData; } + private static ProxyTopicRouteData createProxyTopicRouteDataWithMissingBrokerBeforeValidBroker() { + ProxyTopicRouteData proxyTopicRouteData = createProxyTopicRouteData(2, 2, PermName.PERM_READ | PermName.PERM_WRITE); + QueueData missingBrokerQueueData = createQueueData(2, 2, PermName.PERM_READ | PermName.PERM_WRITE); + missingBrokerQueueData.setBrokerName("missingBrokerName"); + proxyTopicRouteData.getQueueDatas().add(0, missingBrokerQueueData); + return proxyTopicRouteData; + } + @Test public void testGenPartitionFromQueueData() throws Exception { // test queueData with 8 read queues, 8 write queues, and rw permission, expect 8 rw queues.