From 2da929357553acbdd6a0f64932486cd12b4fe9c2 Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Fri, 8 Apr 2022 17:02:03 +0800 Subject: [PATCH] [ISSUE #3949] Refector by code review * Add toString for MetadataHeader * Use RemotingHelper * Use Map in RequestMapping * Add Epoll in GrpcServer --- .../rocketmq/acl/common/MetadataHeader.java | 32 +++--- .../acl/plain/PlainAccessValidator.java | 3 +- .../rocketmq/proxy/config/ProxyConfig.java | 9 ++ .../rocketmq/proxy/grpc/GrpcServer.java | 25 +++-- .../proxy/grpc/v2/adapter/RequestMapping.java | 105 +++++++++--------- 5 files changed, 96 insertions(+), 78 deletions(-) diff --git a/acl/src/main/java/org/apache/rocketmq/acl/common/MetadataHeader.java b/acl/src/main/java/org/apache/rocketmq/acl/common/MetadataHeader.java index a6824918a3..96c7ac7790 100644 --- a/acl/src/main/java/org/apache/rocketmq/acl/common/MetadataHeader.java +++ b/acl/src/main/java/org/apache/rocketmq/acl/common/MetadataHeader.java @@ -122,21 +122,6 @@ public class MetadataHeader { this.datetime, this.sessionToken, this.requestId, this.language, this.clientVersion, this.protocol, this.requestCode); } - - @Override public String toString() { - return "MetadataHeaderBuilder{" + "remoteAddress='" + remoteAddress + '\'' + - ", tenantId='" + tenantId + '\'' + - ", namespace='" + namespace + '\'' + - ", authorization='" + authorization + '\'' + - ", datetime='" + datetime + '\'' + - ", sessionToken='" + sessionToken + '\'' + - ", requestId='" + requestId + '\'' + - ", language='" + language + '\'' + - ", clientVersion='" + clientVersion + '\'' + - ", protocol='" + protocol + '\'' + - ", requestCode=" + requestCode + - '}'; - } } public static MetadataHeader.MetadataHeaderBuilder builder() { @@ -230,4 +215,21 @@ public class MetadataHeader { public void setRequestCode(int requestCode) { this.requestCode = requestCode; } + + @Override public String toString() { + final StringBuilder sb = new StringBuilder("MetadataHeader{"); + sb.append("remoteAddress='").append(remoteAddress).append('\''); + sb.append(", tenantId='").append(tenantId).append('\''); + sb.append(", namespace='").append(namespace).append('\''); + sb.append(", authorization='").append(authorization).append('\''); + sb.append(", datetime='").append(datetime).append('\''); + sb.append(", sessionToken='").append(sessionToken).append('\''); + sb.append(", requestId='").append(requestId).append('\''); + sb.append(", language='").append(language).append('\''); + sb.append(", clientVersion='").append(clientVersion).append('\''); + sb.append(", protocol='").append(protocol).append('\''); + sb.append(", requestCode=").append(requestCode); + sb.append('}'); + return sb.toString(); + } } diff --git a/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java b/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java index 6e1f78463e..cf16b347d5 100644 --- a/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java +++ b/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java @@ -53,6 +53,7 @@ import org.apache.rocketmq.common.protocol.header.UpdateConsumerOffsetRequestHea import org.apache.rocketmq.common.protocol.heartbeat.ConsumerData; import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; +import org.apache.rocketmq.remoting.common.RemotingHelper; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import static org.apache.rocketmq.acl.plain.PlainAccessResource.getRetryTopic; @@ -156,7 +157,7 @@ public class PlainAccessValidator implements AccessValidator { PlainAccessResource accessResource = new PlainAccessResource(); String remoteAddress = header.getRemoteAddress(); if (remoteAddress != null && remoteAddress.contains(":")) { - accessResource.setWhiteRemoteAddress(remoteAddress.substring(0, remoteAddress.lastIndexOf(':'))); + accessResource.setWhiteRemoteAddress(RemotingHelper.parseHostFromAddress(remoteAddress)); } else { accessResource.setWhiteRemoteAddress(remoteAddress); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java index 37b4281430..1375abc3e7 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java @@ -47,6 +47,7 @@ public class ProxyConfig { private String grpcTlsCertPath = ConfigurationManager.getProxyHome() + "/conf/tls/gRPC.chain.cert.pem"; private int grpcBossLoopNum = 1; private int grpcWorkerLoopNum = Runtime.getRuntime().availableProcessors() * 2; + private boolean enableGrpcEpoll = false; private int grpcThreadPoolNums = 16 + Runtime.getRuntime().availableProcessors() * 2; private int grpcThreadPoolQueueCapacity = 100000; private String brokerConfigPath = ConfigurationManager.getProxyHome() + "/conf/broker.conf"; @@ -196,6 +197,14 @@ public class ProxyConfig { this.grpcWorkerLoopNum = grpcWorkerLoopNum; } + public boolean isEnableGrpcEpoll() { + return enableGrpcEpoll; + } + + public void setEnableGrpcEpoll(boolean enableGrpcEpoll) { + this.enableGrpcEpoll = enableGrpcEpoll; + } + public int getGrpcThreadPoolNums() { return grpcThreadPoolNums; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java index 9b22e8d07b..2167ac08cd 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java @@ -19,6 +19,8 @@ package org.apache.rocketmq.proxy.grpc; import io.grpc.netty.shaded.io.grpc.netty.GrpcSslContexts; import io.grpc.netty.shaded.io.grpc.netty.NettyServerBuilder; +import io.grpc.netty.shaded.io.netty.channel.epoll.EpollEventLoopGroup; +import io.grpc.netty.shaded.io.netty.channel.epoll.EpollServerSocketChannel; import io.grpc.netty.shaded.io.netty.channel.nio.NioEventLoopGroup; import io.grpc.netty.shaded.io.netty.channel.socket.nio.NioServerSocketChannel; import io.grpc.netty.shaded.io.netty.handler.ssl.ClientAuth; @@ -39,8 +41,8 @@ import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.interceptor.AuthenticationInterceptor; import org.apache.rocketmq.proxy.grpc.interceptor.ContextInterceptor; import org.apache.rocketmq.proxy.grpc.interceptor.HeaderInterceptor; -import org.apache.rocketmq.proxy.grpc.v2.service.GrpcForwardService; import org.apache.rocketmq.proxy.grpc.v2.GrpcMessagingProcessor; +import org.apache.rocketmq.proxy.grpc.v2.service.GrpcForwardService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -91,12 +93,21 @@ public class GrpcServer implements StartAndShutdown { int workerLoopNum = ConfigurationManager.getProxyConfig().getGrpcWorkerLoopNum(); int maxInboundMessageSize = ConfigurationManager.getProxyConfig().getGrpcMaxInboundMessageSize(); - serverBuilder.maxInboundMessageSize(maxInboundMessageSize) - .bossEventLoopGroup(new NioEventLoopGroup(bossLoopNum)) - .workerEventLoopGroup(new NioEventLoopGroup(workerLoopNum)) - .channelType(NioServerSocketChannel.class) - .addService(messagingProcessor) - .executor(this.executor); + if (ConfigurationManager.getProxyConfig().isEnableGrpcEpoll()) { + serverBuilder.maxInboundMessageSize(maxInboundMessageSize) + .bossEventLoopGroup(new EpollEventLoopGroup(bossLoopNum)) + .workerEventLoopGroup(new EpollEventLoopGroup(workerLoopNum)) + .channelType(EpollServerSocketChannel.class) + .addService(messagingProcessor) + .executor(this.executor); + } else { + serverBuilder.maxInboundMessageSize(maxInboundMessageSize) + .bossEventLoopGroup(new NioEventLoopGroup(bossLoopNum)) + .workerEventLoopGroup(new NioEventLoopGroup(workerLoopNum)) + .channelType(NioServerSocketChannel.class) + .addService(messagingProcessor) + .executor(this.executor); + } // grpc interceptors, including acl, logging etc. if (ConfigurationManager.getProxyConfig().isEnableACL()) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java index fe00bd12a4..2a25ff8eb7 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java @@ -17,65 +17,60 @@ package org.apache.rocketmq.proxy.grpc.v2.adapter; -import apache.rocketmq.v1.AckMessageRequest; -import apache.rocketmq.v1.ChangeInvisibleDurationRequest; -import apache.rocketmq.v1.EndTransactionRequest; -import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueResponse; -import apache.rocketmq.v1.HealthCheckRequest; -import apache.rocketmq.v1.HeartbeatRequest; -import apache.rocketmq.v1.NackMessageRequest; -import apache.rocketmq.v1.NotifyClientTerminationRequest; -import apache.rocketmq.v1.PullMessageRequest; -import apache.rocketmq.v1.QueryAssignmentRequest; -import apache.rocketmq.v1.QueryOffsetRequest; -import apache.rocketmq.v1.QueryRouteRequest; -import apache.rocketmq.v1.ReceiveMessageRequest; -import apache.rocketmq.v1.SendMessageRequest; +import apache.rocketmq.v2.AckMessageRequest; +import apache.rocketmq.v2.ChangeInvisibleDurationRequest; +import apache.rocketmq.v2.EndTransactionRequest; +import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; +import apache.rocketmq.v2.HeartbeatRequest; +import apache.rocketmq.v2.NackMessageRequest; +import apache.rocketmq.v2.NotifyClientTerminationRequest; +import apache.rocketmq.v2.PullMessageRequest; +import apache.rocketmq.v2.QueryAssignmentRequest; +import apache.rocketmq.v2.QueryOffsetRequest; +import apache.rocketmq.v2.QueryRouteRequest; +import apache.rocketmq.v2.ReceiveMessageRequest; +import apache.rocketmq.v2.SendMessageRequest; +import java.util.HashMap; +import java.util.Map; import org.apache.rocketmq.common.protocol.RequestCode; public class RequestMapping { + private final static Map REQUEST_MAP = new HashMap() {{ + // v2 + put(QueryRouteRequest.getDescriptor().getFullName(), RequestCode.GET_ROUTEINFO_BY_TOPIC); + put(HeartbeatRequest.getDescriptor().getFullName(), RequestCode.HEART_BEAT); + put(SendMessageRequest.getDescriptor().getFullName(), RequestCode.SEND_MESSAGE_V2); + put(QueryAssignmentRequest.getDescriptor().getFullName(), RequestCode.GET_ROUTEINFO_BY_TOPIC); + put(ReceiveMessageRequest.getDescriptor().getFullName(), RequestCode.PULL_MESSAGE); + put(AckMessageRequest.getDescriptor().getFullName(), RequestCode.UPDATE_CONSUMER_OFFSET); + put(NackMessageRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + put(ForwardMessageToDeadLetterQueueResponse.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + put(EndTransactionRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + put(QueryOffsetRequest.getDescriptor().getFullName(), RequestCode.SEARCH_OFFSET_BY_TIMESTAMP); + put(PullMessageRequest.getDescriptor().getFullName(), RequestCode.PULL_MESSAGE); + put(NotifyClientTerminationRequest.getDescriptor().getFullName(), RequestCode.UNREGISTER_CLIENT); + put(ChangeInvisibleDurationRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + + // v1 + put(apache.rocketmq.v1.QueryRouteRequest.getDescriptor().getFullName(), RequestCode.GET_ROUTEINFO_BY_TOPIC); + put(apache.rocketmq.v1.HeartbeatRequest.getDescriptor().getFullName(), RequestCode.HEART_BEAT); + put(apache.rocketmq.v1.HealthCheckRequest.getDescriptor().getFullName(), RequestCode.HEART_BEAT); + put(apache.rocketmq.v1.SendMessageRequest.getDescriptor().getFullName(), RequestCode.SEND_MESSAGE_V2); + put(apache.rocketmq.v1.QueryAssignmentRequest.getDescriptor().getFullName(), RequestCode.GET_ROUTEINFO_BY_TOPIC); + put(apache.rocketmq.v1.ReceiveMessageRequest.getDescriptor().getFullName(), RequestCode.PULL_MESSAGE); + put(apache.rocketmq.v1.AckMessageRequest.getDescriptor().getFullName(), RequestCode.UPDATE_CONSUMER_OFFSET); + put(apache.rocketmq.v1.NackMessageRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + put(apache.rocketmq.v1.ForwardMessageToDeadLetterQueueResponse.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + put(apache.rocketmq.v1.EndTransactionRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + put(apache.rocketmq.v1.QueryOffsetRequest.getDescriptor().getFullName(), RequestCode.SEARCH_OFFSET_BY_TIMESTAMP); + put(apache.rocketmq.v1.PullMessageRequest.getDescriptor().getFullName(), RequestCode.PULL_MESSAGE); + put(apache.rocketmq.v1.NotifyClientTerminationRequest.getDescriptor().getFullName(), RequestCode.UNREGISTER_CLIENT); + put(apache.rocketmq.v1.ChangeInvisibleDurationRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); + }}; + public static int map(String rpcFullName) { - if (QueryRouteRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.GET_ROUTEINFO_BY_TOPIC; - } - if (HeartbeatRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.HEART_BEAT; - } - if (HealthCheckRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.HEART_BEAT; - } - if (SendMessageRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.SEND_MESSAGE_V2; - } - if (QueryAssignmentRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.GET_ROUTEINFO_BY_TOPIC; - } - if (ReceiveMessageRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.PULL_MESSAGE; - } - if (AckMessageRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.UPDATE_CONSUMER_OFFSET; - } - if (NackMessageRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.CONSUMER_SEND_MSG_BACK; - } - if (ForwardMessageToDeadLetterQueueResponse.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.CONSUMER_SEND_MSG_BACK; - } - if (EndTransactionRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.END_TRANSACTION; - } - if (QueryOffsetRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.SEARCH_OFFSET_BY_TIMESTAMP; - } - if (PullMessageRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.PULL_MESSAGE; - } - if (NotifyClientTerminationRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.UNREGISTER_CLIENT; - } - if (ChangeInvisibleDurationRequest.getDescriptor().getFullName().equals(rpcFullName)) { - return RequestCode.CONSUMER_SEND_MSG_BACK; + if (REQUEST_MAP.containsKey(rpcFullName)) { + return REQUEST_MAP.get(rpcFullName); } return RequestCode.HEART_BEAT; }