From 7fd29b43d4476784927e67920f3952c2a7d1a648 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Tue, 12 Apr 2022 13:52:49 +0800 Subject: [PATCH] [ISSUE #3949] Do refactor some code for readability. --- .../rocketmq/proxy/HealthCheckServer.java | 10 +-- .../apache/rocketmq/proxy/ProxyStartup.java | 2 +- .../proxy/grpc/v2/adapter/RequestMapping.java | 64 ++++++++++--------- .../grpc/v2/service/ClusterGrpcService.java | 4 +- 4 files changed, 42 insertions(+), 38 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java index d058dccb28..2a6985c611 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java @@ -23,9 +23,9 @@ import com.sun.net.httpserver.HttpServer; import java.io.IOException; import java.io.OutputStream; import java.net.InetSocketAddress; -import java.util.concurrent.TimeUnit; -import org.apache.rocketmq.proxy.config.ConfigurationManager; +import java.time.Duration; import org.apache.rocketmq.proxy.common.StartAndShutdown; +import org.apache.rocketmq.proxy.config.ConfigurationManager; public class HealthCheckServer implements StartAndShutdown { @@ -34,7 +34,8 @@ public class HealthCheckServer implements StartAndShutdown { @Override public void start() throws Exception { this.healthChecker = HttpServer.create( - new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0); + new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0 + ); this.healthChecker.createContext("/status", new HealthCheckHandler()); this.healthChecker.setExecutor(null); this.healthChecker.start(); @@ -43,7 +44,8 @@ public class HealthCheckServer implements StartAndShutdown { @Override public void shutdown() throws InterruptedException { this.healthChecker.stop(0); - Thread.sleep(TimeUnit.SECONDS.toMillis(ConfigurationManager.getProxyConfig().getWaitAfterStopHealthCheckInSeconds())); + long waitAfterStopHealthCheckInSeconds = ConfigurationManager.getProxyConfig().getWaitAfterStopHealthCheckInSeconds(); + Thread.sleep(Duration.ofSeconds(waitAfterStopHealthCheckInSeconds).toMillis()); } static class HealthCheckHandler implements HttpHandler { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java index accb214777..57239a7954 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -73,7 +73,7 @@ public class ProxyStartup { try { PROXY_START_AND_SHUTDOWN.shutdown(); } catch (Exception e) { - log.error("err when shutdown proxy", e); + log.error("err when shutdown rmq-proxy", e); } })); } catch (Exception e) { 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 2a25ff8eb7..add450e5a1 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 @@ -35,38 +35,40 @@ 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); + 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); - }}; + // 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 (REQUEST_MAP.containsKey(rpcFullName)) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java index e1081ddd02..b44d958a50 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java @@ -72,7 +72,8 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME); private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor( - new ThreadFactoryImpl("ClusterGrpcServiceScheduledThread")); + new ThreadFactoryImpl("ClusterGrpcServiceScheduledThread") + ); private final ChannelManager channelManager; private final ConnectorManager connectorManager; @@ -179,7 +180,6 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc @Override public void start() throws Exception { - } @Override