From b425608b1bf03a60a5ce0d565c8aecc45ba5cbc4 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Thu, 17 Mar 2022 20:08:37 +0800 Subject: [PATCH] [ISSUE #3949] Refactor partition generation and add unit test. --- .../grpc/service/ClusterGrpcService.java | 52 +--- .../grpc/service/cluster/RouteService.java | 58 ++-- .../service/cluster/RouteServiceTest.java | 257 +++++++----------- 3 files changed, 130 insertions(+), 237 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 8db0a48d63..be0368f9b4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -17,47 +17,9 @@ package org.apache.rocketmq.proxy.grpc.service; -import apache.rocketmq.v1.AckMessageRequest; -import apache.rocketmq.v1.AckMessageResponse; -import apache.rocketmq.v1.ChangeInvisibleDurationRequest; -import apache.rocketmq.v1.ChangeInvisibleDurationResponse; -import apache.rocketmq.v1.EndTransactionRequest; -import apache.rocketmq.v1.EndTransactionResponse; -import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueRequest; -import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueResponse; -import apache.rocketmq.v1.HealthCheckRequest; -import apache.rocketmq.v1.HealthCheckResponse; -import apache.rocketmq.v1.HeartbeatRequest; -import apache.rocketmq.v1.HeartbeatResponse; -import apache.rocketmq.v1.NackMessageRequest; -import apache.rocketmq.v1.NackMessageResponse; -import apache.rocketmq.v1.NoopCommand; -import apache.rocketmq.v1.NotifyClientTerminationRequest; -import apache.rocketmq.v1.NotifyClientTerminationResponse; -import apache.rocketmq.v1.PollCommandRequest; -import apache.rocketmq.v1.PollCommandResponse; -import apache.rocketmq.v1.PullMessageRequest; -import apache.rocketmq.v1.PullMessageResponse; -import apache.rocketmq.v1.QueryAssignmentRequest; -import apache.rocketmq.v1.QueryAssignmentResponse; -import apache.rocketmq.v1.QueryOffsetRequest; -import apache.rocketmq.v1.QueryOffsetResponse; -import apache.rocketmq.v1.QueryRouteRequest; -import apache.rocketmq.v1.QueryRouteResponse; -import apache.rocketmq.v1.ReceiveMessageRequest; -import apache.rocketmq.v1.ReceiveMessageResponse; -import apache.rocketmq.v1.ReportMessageConsumptionResultRequest; -import apache.rocketmq.v1.ReportMessageConsumptionResultResponse; -import apache.rocketmq.v1.ReportThreadStackTraceRequest; -import apache.rocketmq.v1.ReportThreadStackTraceResponse; -import apache.rocketmq.v1.Resource; -import apache.rocketmq.v1.SendMessageRequest; -import apache.rocketmq.v1.SendMessageResponse; +import apache.rocketmq.v1.*; import com.google.rpc.Code; import io.grpc.Context; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; import org.apache.rocketmq.common.ThreadFactoryImpl; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.proxy.channel.ChannelManager; @@ -69,15 +31,14 @@ import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.common.Converter; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService; -import org.apache.rocketmq.proxy.grpc.service.cluster.PullMessageService; -import org.apache.rocketmq.proxy.grpc.service.cluster.ReceiveMessageService; -import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService; -import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService; -import org.apache.rocketmq.proxy.grpc.service.cluster.TransactionService; +import org.apache.rocketmq.proxy.grpc.service.cluster.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; + public class ClusterGrpcService extends AbstractStartAndShutdown implements GrpcForwardService { private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME); @@ -85,7 +46,6 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc new ThreadFactoryImpl("ClusterGrpcServiceScheduledThread")); private final ChannelManager channelManager; - private final ConnectorManager connectorManager; private final ProducerService producerService; private final ReceiveMessageService receiveMessageService; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java index 2b4a34fc3b..ea8ad005fe 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java @@ -16,33 +16,25 @@ */ package org.apache.rocketmq.proxy.grpc.service.cluster; -import apache.rocketmq.v1.Assignment; -import apache.rocketmq.v1.Broker; -import apache.rocketmq.v1.Endpoints; -import apache.rocketmq.v1.Partition; -import apache.rocketmq.v1.Permission; -import apache.rocketmq.v1.QueryAssignmentRequest; -import apache.rocketmq.v1.QueryAssignmentResponse; -import apache.rocketmq.v1.QueryRouteRequest; -import apache.rocketmq.v1.QueryRouteResponse; -import apache.rocketmq.v1.Resource; +import apache.rocketmq.v1.*; import com.google.rpc.Code; import io.grpc.Context; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.constant.PermName; import org.apache.rocketmq.common.protocol.route.QueueData; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.proxy.connector.ConnectorManager; -import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; 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.common.Converter; import org.apache.rocketmq.proxy.grpc.common.ParameterConverter; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; + public class RouteService extends BaseService { private volatile ParameterConverter queryRouteEndpointConverter; @@ -148,21 +140,33 @@ public class RouteService extends BaseService { r = queueData.getReadQueueNums(); } - for (int i = 0; i < (rw + r + w); i++) { - Partition.Builder builder = Partition.newBuilder() + // r here means readOnly queue nums, w means writeOnly queue nums, while rw means readable and writable queue nums. + int queueIdIndex = 0; + for(int i = 0; i < r; i++){ + Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.READ); + partitionList.add(partition); + } + + for(int i = 0; i < w; i++){ + Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.WRITE); + partitionList.add(partition); + } + + for (int i = 0; i < rw; i++) { + Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.READ_WRITE); + partitionList.add(partition); + } + + return partitionList; + } + + private static Partition buildPartition(Broker broker, Resource topic, int queueId, Permission perm) { + Partition.Builder builder = Partition.newBuilder() .setBroker(broker) .setTopic(topic) - .setId(i); - if (i < r) { - builder.setPermission(Permission.READ); - } else if (i < w) { - builder.setPermission(Permission.WRITE); - } else { - builder.setPermission(Permission.READ_WRITE); - } - partitionList.add(builder.build()); - } - return partitionList; + .setId(queueId); + builder.setPermission(perm); + return builder.build(); } public CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java index 03a0439302..41851d9b8b 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java @@ -1,177 +1,106 @@ +/* + * 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.service.cluster; -import apache.rocketmq.v1.Address; -import apache.rocketmq.v1.AddressScheme; -import apache.rocketmq.v1.Endpoints; -import apache.rocketmq.v1.QueryAssignmentRequest; -import apache.rocketmq.v1.QueryAssignmentResponse; -import apache.rocketmq.v1.QueryRouteRequest; -import apache.rocketmq.v1.QueryRouteResponse; +import apache.rocketmq.v1.Broker; +import apache.rocketmq.v1.Partition; +import apache.rocketmq.v1.Permission; import apache.rocketmq.v1.Resource; -import com.google.rpc.Code; -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.constant.PermName; 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.common.ResponseBuilder; +import org.junit.After; +import org.junit.Before; import org.junit.Test; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNull; -import static org.mockito.Mockito.when; +import java.util.List; -public class RouteServiceTest extends BaseServiceTest { +import static org.assertj.core.api.Assertions.assertThat; - @Override - public void beforeEach() throws Throwable { - TopicRouteData routeData = new TopicRouteData(); +public class RouteServiceTest { + public static final String BROKER_NAME = "brokerName"; + public static final String NAMESPACE = "namespace"; + public static final String TOPIC = "topic"; + public static final Broker MOCK_BROKER = Broker.newBuilder().setName(BROKER_NAME).build(); + public static final Resource MOCK_TOPIC = Resource.newBuilder() + .setName(TOPIC) + .setResourceNamespace(NAMESPACE) + .build(); - List brokerDataList = new ArrayList<>(); - BrokerData brokerData = new BrokerData(); - brokerData.setCluster("cluster"); - brokerData.setBrokerName("brokerName"); - HashMap brokerAddrs = new HashMap() {{ - put(0L, "127.0.0.1:10911"); - }}; - brokerData.setBrokerAddrs(brokerAddrs); - brokerDataList.add(brokerData); + @Before + public void before() throws Exception { + } - List queueDataList = new ArrayList<>(); + @After + public void after() throws Exception { + } + + + @Test + public void testGenPartitionFromQueueData() throws Exception { + // test queueData with 8 read queues, 8 write queues, and rw permission, expect 8 rw queues. + QueueData queueDataWith8R8WPermRW = mockQueueData(8, 8, PermName.PERM_READ | PermName.PERM_WRITE); + List partitionWith8R8WPermRW = RouteService.genPartitionFromQueueData(queueDataWith8R8WPermRW, MOCK_TOPIC, MOCK_BROKER); + assertThat(partitionWith8R8WPermRW.size()).isEqualTo(8); + assertThat(partitionWith8R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(8); + assertThat(partitionWith8R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(0); + assertThat(partitionWith8R8WPermRW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(0); + + // test queueData with 8 read queues, 8 write queues, and read only permission, expect 8 read only queues. + QueueData queueDataWith8R8WPermR = mockQueueData(8, 8, PermName.PERM_READ); + List partitionWith8R8WPermR = RouteService.genPartitionFromQueueData(queueDataWith8R8WPermR, MOCK_TOPIC, MOCK_BROKER); + assertThat(partitionWith8R8WPermR.size()).isEqualTo(8); + assertThat(partitionWith8R8WPermR.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(8); + assertThat(partitionWith8R8WPermR.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(0); + assertThat(partitionWith8R8WPermR.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(0); + + // test queueData with 8 read queues, 8 write queues, and write only permission, expect 8 write only queues. + QueueData queueDataWith8R8WPermW = mockQueueData(8, 8, PermName.PERM_WRITE); + List partitionWith8R8WPermW = RouteService.genPartitionFromQueueData(queueDataWith8R8WPermW, MOCK_TOPIC, MOCK_BROKER); + assertThat(partitionWith8R8WPermW.size()).isEqualTo(8); + assertThat(partitionWith8R8WPermW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(8); + assertThat(partitionWith8R8WPermW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(0); + assertThat(partitionWith8R8WPermW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(0); + + // test queueData with 8 read queues, 0 write queues, and rw permission, expect 8 read only queues. + QueueData queueDataWith8R0WPermRW = mockQueueData(8, 0, PermName.PERM_READ | PermName.PERM_WRITE); + List partitionWith8R0WPermRW = RouteService.genPartitionFromQueueData(queueDataWith8R0WPermRW, MOCK_TOPIC, MOCK_BROKER); + assertThat(partitionWith8R0WPermRW.size()).isEqualTo(8); + assertThat(partitionWith8R0WPermRW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(8); + assertThat(partitionWith8R0WPermRW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(0); + assertThat(partitionWith8R0WPermRW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(0); + + // test queueData with 4 read queues, 8 write queues, and rw permission, expect 4 rw queues and 4 write only queues. + QueueData queueDataWith4R8WPermRW = mockQueueData(4, 8, PermName.PERM_READ | PermName.PERM_WRITE); + List partitionWith4R8WPermRW = RouteService.genPartitionFromQueueData(queueDataWith4R8WPermRW, MOCK_TOPIC, MOCK_BROKER); + assertThat(partitionWith4R8WPermRW.size()).isEqualTo(8); + assertThat(partitionWith4R8WPermRW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(4); + assertThat(partitionWith4R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(4); + assertThat(partitionWith4R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(0); + + } + + private QueueData mockQueueData(int r, int w, int perm) { 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, "")); + queueData.setBrokerName(BROKER_NAME); + queueData.setReadQueueNums(r); + queueData.setWriteQueueNums(w); + queueData.setPerm(perm); + return queueData; } - @Test - public void testQueryRouteWithInvalidEndpoints() { - RouteService routeService = new RouteService(this.clientManager); - - CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() - .build()); - - try { - QueryRouteResponse response = future.get(); - assertEquals(Code.INVALID_ARGUMENT.getNumber(), response.getCommon().getStatus().getCode()); - } catch (Exception e) { - assertNull(e); - } - } - - @Test - public void testQueryRoute() { - RouteService routeService = new RouteService(this.clientManager); - - CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() - .setEndpoints(Endpoints.newBuilder() - .addAddresses(Address.newBuilder() - .setPort(80) - .setHost("host") - .build()) - .setScheme(AddressScheme.DOMAIN_NAME) - .build()) - .setTopic(Resource.newBuilder() - .setName("topic") - .build()) - .build()); - - try { - QueryRouteResponse response = future.get(); - assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode()); - assertEquals(8, response.getPartitionsCount()); - assertEquals("host", response.getPartitions(0).getBroker() - .getEndpoints().getAddresses(0).getHost()); - } catch (Exception e) { - assertNull(e); - } - } - - @Test - public void testQueryRouteWhenTopicNotExist() { - RouteService routeService = new RouteService(this.clientManager); - - CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() - .setEndpoints(Endpoints.newBuilder() - .addAddresses(Address.newBuilder() - .setPort(80) - .setHost("host") - .build()) - .setScheme(AddressScheme.DOMAIN_NAME) - .build()) - .setTopic(Resource.newBuilder() - .setName("notExistTopic") - .build()) - .build()); - - try { - QueryRouteResponse response = future.get(); - assertEquals(Code.NOT_FOUND.getNumber(), response.getCommon().getStatus().getCode()); - } catch (Exception e) { - assertNull(e); - } - } - - @Test - public void testQueryAssignmentInvalidEndpoints() { - RouteService routeService = new RouteService(this.clientManager); - - CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() - .build()); - - try { - QueryAssignmentResponse response = future.get(); - assertEquals(Code.INVALID_ARGUMENT.getNumber(), response.getCommon().getStatus().getCode()); - } catch (Exception e) { - assertNull(e); - } - } - - @Test - public void testQueryAssignment() { - RouteService routeService = new RouteService(this.clientManager); - - CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() - .setEndpoints(Endpoints.newBuilder() - .addAddresses(Address.newBuilder() - .setPort(80) - .setHost("host") - .build()) - .setScheme(AddressScheme.DOMAIN_NAME) - .build()) - .setTopic(Resource.newBuilder() - .setName("topic") - .build()) - .setGroup(Resource.newBuilder() - .setName("group") - .build()) - .setClientId("clientId") - .build()); - - try { - QueryAssignmentResponse response = future.get(); - assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode()); - assertEquals(1, response.getAssignmentsCount()); - assertEquals("brokerName", response.getAssignments(0).getPartition().getBroker().getName()); - assertEquals("host", response.getAssignments(0).getPartition().getBroker().getEndpoints().getAddresses(0).getHost()); - } catch (Exception e) { - assertNull(e); - } - } -} \ No newline at end of file +}