diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/AbstractRouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/AbstractRouteService.java new file mode 100644 index 0000000000..d582d718a2 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/AbstractRouteService.java @@ -0,0 +1,100 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.proxy.grpc.v2.service; + +import apache.rocketmq.v2.Endpoints; +import apache.rocketmq.v2.QueryAssignmentRequest; +import apache.rocketmq.v2.QueryAssignmentResponse; +import apache.rocketmq.v2.QueryRouteRequest; +import apache.rocketmq.v2.QueryRouteResponse; +import io.grpc.Context; +import java.util.concurrent.CompletableFuture; +import org.apache.rocketmq.proxy.common.ParameterConverter; +import org.apache.rocketmq.proxy.connector.ConnectorManager; +import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook; +import org.apache.rocketmq.proxy.grpc.v2.service.cluster.AssignmentQueueSelector; +import org.apache.rocketmq.proxy.grpc.v2.service.cluster.BaseService; +import org.apache.rocketmq.proxy.grpc.v2.service.cluster.DefaultAssignmentQueueSelector; + +public abstract class AbstractRouteService extends BaseService { + protected volatile ParameterConverter queryRouteEndpointConverter; + protected volatile ResponseHook queryRouteHook; + + protected volatile ParameterConverter queryAssignmentEndpointConverter; + protected volatile ResponseHook queryAssignmentHook; + protected volatile AssignmentQueueSelector assignmentQueueSelector; + + protected final GrpcClientManager grpcClientManager; + + public AbstractRouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) { + super(connectorManager); + this.grpcClientManager = grpcClientManager; + this.queryRouteEndpointConverter = (ctx, parameter) -> parameter; + this.queryAssignmentEndpointConverter = (ctx, parameter) -> parameter; + this.assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache()); + } + + public abstract CompletableFuture queryRoute(Context ctx, QueryRouteRequest request); + + public abstract CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request); + + public ParameterConverter getQueryRouteEndpointConverter() { + return queryRouteEndpointConverter; + } + + public void setQueryRouteEndpointConverter( + ParameterConverter queryRouteEndpointConverter) { + this.queryRouteEndpointConverter = queryRouteEndpointConverter; + } + + public ResponseHook getQueryRouteHook() { + return queryRouteHook; + } + + public void setQueryRouteHook( + ResponseHook queryRouteHook) { + this.queryRouteHook = queryRouteHook; + } + + public ParameterConverter getQueryAssignmentEndpointConverter() { + return queryAssignmentEndpointConverter; + } + + public void setQueryAssignmentEndpointConverter( + ParameterConverter queryAssignmentEndpointConverter) { + this.queryAssignmentEndpointConverter = queryAssignmentEndpointConverter; + } + + public AssignmentQueueSelector getAssignmentQueueSelector() { + return assignmentQueueSelector; + } + + public void setAssignmentQueueSelector( + AssignmentQueueSelector assignmentQueueSelector) { + this.assignmentQueueSelector = assignmentQueueSelector; + } + + public ResponseHook getQueryAssignmentHook() { + return queryAssignmentHook; + } + + public void setQueryAssignmentHook( + ResponseHook queryAssignmentHook) { + this.queryAssignmentHook = queryAssignmentHook; + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java index 920c5b2938..070c1e7a4c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java @@ -98,7 +98,6 @@ import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException; -import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseWriter; import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel; @@ -106,7 +105,7 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.ReceiveMessageChannel; import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.SendMessageChannel; import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.ReceiveMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.SendMessageResponseHandler; -import org.apache.rocketmq.proxy.grpc.v2.service.cluster.RouteService; +import org.apache.rocketmq.proxy.grpc.v2.service.local.RouteService; import org.apache.rocketmq.proxy.grpc.v2.service.local.LocalWriteQueueSelector; import org.apache.rocketmq.remoting.RemotingServer; import org.apache.rocketmq.remoting.netty.NettyRemotingAbstract; @@ -143,7 +142,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo ConnectorManager connectorManager = new ConnectorManager(null); this.telemetryCommandManager = telemetryCommandManager; this.grpcClientManager = new GrpcClientManager(); - this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager, grpcClientManager); + this.routeService = new RouteService(connectorManager, grpcClientManager); this.clientSettingsService = new ClientSettingsService(this.channelManager, this.grpcClientManager, this.telemetryCommandManager); this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel()); this.localWriteQueueSelector = new LocalWriteQueueSelector(brokerController.getBrokerConfig().getBrokerName(), diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java index d91bcb4ccc..e8c206b609 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java @@ -33,38 +33,21 @@ import java.util.List; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.protocol.route.QueueData; import org.apache.rocketmq.common.protocol.route.TopicRouteData; -import org.apache.rocketmq.proxy.common.ParameterConverter; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook; +import org.apache.rocketmq.proxy.grpc.v2.service.AbstractRouteService; import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; -public class RouteService extends BaseService { - private volatile ParameterConverter queryRouteEndpointConverter; - private volatile ResponseHook queryRouteHook; - - private volatile ParameterConverter queryAssignmentEndpointConverter; - private volatile AssignmentQueueSelector assignmentQueueSelector; - private volatile ResponseHook queryAssignmentHook; - - protected final GrpcClientManager grpcClientManager; - +public class RouteService extends AbstractRouteService { public RouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) { - super(connectorManager); - this.grpcClientManager = grpcClientManager; + super(connectorManager, grpcClientManager); } @Override - public void start() throws Exception { - this.queryRouteEndpointConverter = (ctx, parameter) -> parameter; - this.queryAssignmentEndpointConverter = (ctx, parameter) -> parameter; - this.assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache()); - } - public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { CompletableFuture future = new CompletableFuture<>(); future.whenComplete((response, throwable) -> { @@ -115,6 +98,7 @@ public class RouteService extends BaseService { return future; } + @Override public CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { CompletableFuture future = new CompletableFuture<>(); future.whenComplete((response, throwable) -> { @@ -164,49 +148,4 @@ public class RouteService extends BaseService { } return future; } - - public ParameterConverter getQueryRouteEndpointConverter() { - return queryRouteEndpointConverter; - } - - public void setQueryRouteEndpointConverter( - ParameterConverter queryRouteEndpointConverter) { - this.queryRouteEndpointConverter = queryRouteEndpointConverter; - } - - public ResponseHook getQueryRouteHook() { - return queryRouteHook; - } - - public void setQueryRouteHook( - ResponseHook queryRouteHook) { - this.queryRouteHook = queryRouteHook; - } - - public ParameterConverter getQueryAssignmentEndpointConverter() { - return queryAssignmentEndpointConverter; - } - - public void setQueryAssignmentEndpointConverter( - ParameterConverter queryAssignmentEndpointConverter) { - this.queryAssignmentEndpointConverter = queryAssignmentEndpointConverter; - } - - public AssignmentQueueSelector getAssignmentQueueSelector() { - return assignmentQueueSelector; - } - - public void setAssignmentQueueSelector( - AssignmentQueueSelector assignmentQueueSelector) { - this.assignmentQueueSelector = assignmentQueueSelector; - } - - public ResponseHook getQueryAssignmentHook() { - return queryAssignmentHook; - } - - public void setQueryAssignmentHook( - ResponseHook queryAssignmentHook) { - this.queryAssignmentHook = queryAssignmentHook; - } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/RouteService.java new file mode 100644 index 0000000000..7083f0684a --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/RouteService.java @@ -0,0 +1,175 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.proxy.grpc.v2.service.local; + +import apache.rocketmq.v2.Address; +import apache.rocketmq.v2.AddressScheme; +import apache.rocketmq.v2.Assignment; +import apache.rocketmq.v2.Broker; +import apache.rocketmq.v2.Code; +import apache.rocketmq.v2.Endpoints; +import apache.rocketmq.v2.MessageQueue; +import apache.rocketmq.v2.Permission; +import apache.rocketmq.v2.QueryAssignmentRequest; +import apache.rocketmq.v2.QueryAssignmentResponse; +import apache.rocketmq.v2.QueryRouteRequest; +import apache.rocketmq.v2.QueryRouteResponse; +import com.google.common.net.HostAndPort; +import io.grpc.Context; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.protocol.route.QueueData; +import org.apache.rocketmq.common.protocol.route.TopicRouteData; +import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.connector.ConnectorManager; +import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; +import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; +import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper; +import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.v2.service.AbstractRouteService; +import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; + +public class RouteService extends AbstractRouteService { + public RouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) { + super(connectorManager, grpcClientManager); + } + + @Override + public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { + CompletableFuture future = new CompletableFuture<>(); + future.whenComplete((response, throwable) -> { + if (queryRouteHook != null) { + queryRouteHook.beforeResponse(ctx, request, response, throwable); + } + }); + + try { + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); + MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName); + TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); + List queueDataList = topicRouteData.getQueueDatas(); + List brokerDataList = topicRouteData.getBrokerDatas(); + + List messageQueueList = new ArrayList<>(); + Map> brokerMap = buildBrokerMap(brokerDataList); + + for (QueueData queueData : queueDataList) { + String brokerName = queueData.getBrokerName(); + Map brokerIdMap = brokerMap.get(brokerName); + if (brokerIdMap == null) { + break; + } + for (Broker broker : brokerIdMap.values()) { + messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker)); + } + } + + QueryRouteResponse response = QueryRouteResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) + .addAllMessageQueues(messageQueueList) + .build(); + future.complete(response); + } catch (Throwable t) { + if (TopicRouteHelper.isTopicNotExistError(t)) { + future.complete(QueryRouteResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.TOPIC_NOT_FOUND, t.getMessage())) + .build()); + } else { + future.completeExceptionally(t); + } + } + return future; + } + + @Override + public CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { + CompletableFuture future = new CompletableFuture<>(); + future.whenComplete((response, throwable) -> { + if (queryAssignmentHook != null) { + queryAssignmentHook.beforeResponse(ctx, request, response, throwable); + } + }); + + try { + List assignments = new ArrayList<>(); + List messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); + MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName); + TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); + Map> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas()); + for (SelectableMessageQueue messageQueue : messageQueueList) { + Map brokerIdMap = brokerMap.get(messageQueue.getBrokerName()); + if (brokerIdMap != null) { + Broker broker = brokerIdMap.get(0L); + + MessageQueue defaultMessageQueue = MessageQueue.newBuilder() + .setTopic(request.getTopic()) + .setId(-1) + .setPermission(Permission.READ_WRITE) + .setBroker(broker) + .build(); + + assignments.add(Assignment.newBuilder() + .setMessageQueue(defaultMessageQueue) + .build()); + } + } + QueryAssignmentResponse response = QueryAssignmentResponse.newBuilder() + .addAllAssignments(assignments) + .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) + .build(); + future.complete(response); + } catch (Throwable t) { + future.completeExceptionally(t); + } + return future; + } + + private Map> buildBrokerMap(List brokerDataList) { + Map> brokerMap = new HashMap<>(); + for (BrokerData brokerData : brokerDataList) { + Map brokerIdMap = new HashMap<>(); + String brokerName = brokerData.getBrokerName(); + for (Map.Entry entry : brokerData.getBrokerAddrs().entrySet()) { + Long brokerId = entry.getKey(); + HostAndPort hostAndPort = HostAndPort.fromString(entry.getValue()); + Broker broker = Broker.newBuilder() + .setName(brokerName) + .setId(Math.toIntExact(brokerId)) + .setEndpoints(Endpoints.newBuilder() + .setScheme(AddressScheme.IPv4) + .addAddresses( + Address.newBuilder() + .setPort(ConfigurationManager.getProxyConfig().getGrpcServerPort()) + .setHost(hostAndPort.getHost()) + ) + .build()) + .build(); + + brokerIdMap.put(brokerId, broker); + } + brokerMap.put(brokerName, brokerIdMap); + } + return brokerMap; + } +} diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/local/RouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/local/RouteServiceTest.java new file mode 100644 index 0000000000..25078d4a89 --- /dev/null +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/local/RouteServiceTest.java @@ -0,0 +1,132 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.proxy.grpc.v2.service.local; + +import apache.rocketmq.v2.Address; +import apache.rocketmq.v2.AddressScheme; +import apache.rocketmq.v2.Code; +import apache.rocketmq.v2.Endpoints; +import apache.rocketmq.v2.QueryAssignmentRequest; +import apache.rocketmq.v2.QueryAssignmentResponse; +import apache.rocketmq.v2.QueryRouteRequest; +import apache.rocketmq.v2.QueryRouteResponse; +import apache.rocketmq.v2.Resource; +import apache.rocketmq.v2.Settings; +import com.google.common.net.HostAndPort; +import io.grpc.Context; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import org.apache.rocketmq.client.exception.MQClientException; +import org.apache.rocketmq.common.protocol.ResponseCode; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.protocol.route.QueueData; +import org.apache.rocketmq.common.protocol.route.TopicRouteData; +import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; +import org.apache.rocketmq.proxy.grpc.v2.service.cluster.BaseServiceTest; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.when; + +public class RouteServiceTest extends BaseServiceTest { + private String brokerAddress = "127.0.0.1:10911"; + private static final Settings WITH_HOST_SETTINGS = Settings.newBuilder() + .setAccessPoint(Endpoints.newBuilder() + .addAddresses(Address.newBuilder() + .setPort(80) + .setHost("host") + .build()) + .setScheme(AddressScheme.DOMAIN_NAME) + .build()) + .build(); + + @Test + public void testLocalModeQueryRoute() throws Exception { + RouteService routeService = new RouteService(this.connectorManager, this.grpcClientManager); + routeService.start(); + + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS); + + CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() + .setTopic(Resource.newBuilder() + .setName("topic") + .build()) + .build()); + QueryRouteResponse response = future.get(); + assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber()); + assertEquals(8, response.getMessageQueuesCount()); + assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getMessageQueues(0).getBroker() + .getEndpoints().getAddresses(0).getHost()); + } + + @Test + public void testLocalModeQueryAssignment() throws Exception { + RouteService routeService = new RouteService(this.connectorManager, this.grpcClientManager); + routeService.start(); + + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS); + + CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() + .setTopic(Resource.newBuilder() + .setName("topic") + .build()) + .setGroup(Resource.newBuilder() + .setName("group") + .build()) + .build()); + + QueryAssignmentResponse response = future.get(); + assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber()); + assertEquals(1, response.getAssignmentsCount()); + assertEquals("brokerName", response.getAssignments(0).getMessageQueue().getBroker().getName()); + assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getAssignments(0).getMessageQueue().getBroker().getEndpoints().getAddresses(0).getHost()); + } + + @Override public void beforeEach() throws Throwable { + TopicRouteData routeData = new TopicRouteData(); + + List brokerDataList = new ArrayList<>(); + BrokerData brokerData = new BrokerData(); + brokerData.setCluster("cluster"); + brokerData.setBrokerName("brokerName"); + HashMap brokerAddrs = new HashMap() {{ + put(0L, brokerAddress); + }}; + brokerData.setBrokerAddrs(brokerAddrs); + brokerDataList.add(brokerData); + + List queueDataList = new ArrayList<>(); + QueueData queueData = new QueueData(); + queueData.setPerm(6); + queueData.setWriteQueueNums(8); + queueData.setReadQueueNums(8); + queueData.setBrokerName("brokerName"); + queueDataList.add(queueData); + + routeData.setBrokerDatas(brokerDataList); + routeData.setQueueDatas(queueDataList); + + MessageQueueWrapper messageQueueWrapper = new MessageQueueWrapper("topic", routeData); + when(this.topicRouteCache.getMessageQueue("topic")).thenReturn(messageQueueWrapper); + + when(this.topicRouteCache.getMessageQueue("notExistTopic")).thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST, "")); + } +} \ No newline at end of file