mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Refector by code review
* Add toString for MetadataHeader * Use RemotingHelper * Use Map in RequestMapping * Add Epoll in GrpcServer
This commit is contained in:
@@ -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();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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()) {
|
||||
|
||||
@@ -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<String, Integer> REQUEST_MAP = new HashMap<String, Integer>() {{
|
||||
// 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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user