From caaadca11c94b587a09a7cac3395ec806bb50cf9 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Fri, 18 Mar 2022 19:17:18 +0800 Subject: [PATCH] [ISSUE #3949] Improve readability. --- .../rocketmq/proxy/grpc/common/Converter.java | 20 +++++++++---------- .../grpc/service/cluster/ClientService.java | 17 ++++++++++------ .../DefaultAssignmentQueueSelector.java | 3 ++- 3 files changed, 23 insertions(+), 17 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java index d6cab5b5c3..7b0e978e60 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java @@ -503,10 +503,8 @@ public class Converter { public static Message buildMessage(MessageExt messageExt) { Map userAttributes = buildUserAttributes(messageExt); SystemAttribute systemAttributes = buildSystemAttributes(messageExt); - Resource topic = Resource.newBuilder() - .setResourceNamespace(NamespaceUtil.getNamespaceFromResource(messageExt.getTopic())) - .setName(NamespaceUtil.withoutNamespace(messageExt.getTopic())) - .build(); + Resource topic = buildResource(messageExt.getTopic()); + return Message.newBuilder() .setTopic(topic) .putAllUserAttribute(userAttributes) @@ -641,12 +639,7 @@ public class Converter { // publisher_group String producerGroup = messageExt.getProperty(MessageConst.PROPERTY_PRODUCER_GROUP); if (producerGroup != null) { - String namespaceId = NamespaceUtil.getNamespaceFromResource(producerGroup); - String group = NamespaceUtil.withoutNamespace(producerGroup); - systemAttributeBuilder.setProducerGroup(Resource.newBuilder() - .setResourceNamespace(namespaceId) - .setName(group) - .build()); + systemAttributeBuilder.setProducerGroup(buildResource(producerGroup)); } // trace context @@ -690,6 +683,13 @@ public class Converter { return consumeMessageDirectlyResult; } + public static Resource buildResource(String resourceNameWithNamespace) { + return Resource.newBuilder() + .setResourceNamespace(NamespaceUtil.getNamespaceFromResource(resourceNameWithNamespace)) + .setName(NamespaceUtil.withoutNamespace(resourceNameWithNamespace)) + .build(); + } + public static UnregisterClientRequestHeader buildUnregisterClientRequestHeader(NotifyClientTerminationRequest request) { UnregisterClientRequestHeader header = new UnregisterClientRequestHeader(); header.setClientID(request.getClientId()); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java index f2f20c517c..9872bd35c3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java @@ -24,6 +24,9 @@ import apache.rocketmq.v1.PollCommandRequest; import apache.rocketmq.v1.PollCommandResponse; import apache.rocketmq.v1.Resource; import io.grpc.Context; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import org.apache.rocketmq.broker.client.ClientChannelInfo; import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; @@ -32,16 +35,12 @@ import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; +import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.remoting.protocol.LanguageCode; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; - public class ClientService extends BaseService { private static final Logger log = LoggerFactory.getLogger(ClientService.class); @@ -52,7 +51,12 @@ public class ClientService extends BaseService { private final ProducerManager producerManager; private final PollCommandResponseManager pollCommandResponseManager; - public ClientService(ConnectorManager connectorManager, ScheduledExecutorService scheduledExecutorService, ChannelManager channelManager, PollCommandResponseManager pollCommandResponseManager) { + public ClientService( + ConnectorManager connectorManager, + ScheduledExecutorService scheduledExecutorService, + ChannelManager channelManager, + PollCommandResponseManager pollCommandResponseManager + ) { super(connectorManager); scheduledExecutorService.scheduleWithFixedDelay(this::scanNotActiveChannel, 1000 * 10, 1000 * 10, TimeUnit.MILLISECONDS); this.channelManager = channelManager; @@ -70,6 +74,7 @@ public class ClientService extends BaseService { if (request.hasProducerData()) { String producerGroup = Converter.getResourceNameWithNamespace(request.getProducerData().getGroup()); GrpcClientChannel channel = GrpcClientChannel.create(channelManager, producerGroup, clientId, pollCommandResponseManager); + //TODO: Use the faked MQ Version ? ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal()); producerManager.registerProducer(producerGroup, clientChannelInfo); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java index 1738cb1f3e..92ea26eb5a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java @@ -34,7 +34,8 @@ public class DefaultAssignmentQueueSelector implements AssignmentQueueSelector { @Override public List getAssignment(Context ctx, QueryAssignmentRequest request) throws Exception { - MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic())); + String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(topicName); return messageQueueWrapper.getReadSelector().getBrokerActingQueues(); } } \ No newline at end of file