mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Do refactor some code for readability.
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -35,38 +35,40 @@ 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);
|
||||
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);
|
||||
}};
|
||||
// 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)) {
|
||||
|
||||
+2
-2
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user