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 789a0dcf29..d058dccb28 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/HealthCheckServer.java @@ -33,7 +33,8 @@ public class HealthCheckServer implements StartAndShutdown { @Override public void start() throws Exception { - this.healthChecker = HttpServer.create(new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0); + this.healthChecker = HttpServer.create( + new InetSocketAddress(ConfigurationManager.getProxyConfig().getHealthCheckPort()), 0); this.healthChecker.createContext("/status", new HealthCheckHandler()); this.healthChecker.setExecutor(null); this.healthChecker.start(); 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 810fd8c1ed..6648992efc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -39,11 +39,12 @@ import org.slf4j.LoggerFactory; public class ProxyStartup { - private static final Logger log = LoggerFactory.getLogger(ProxyStartup.class); - private static final ProxyStartAndShutdown proxyStartAndShutdown = new ProxyStartAndShutdown(); + private static final Logger LOGGER = LoggerFactory.getLogger(ProxyStartup.class); + private static final ProxyStartAndShutdown PROXY_START_AND_SHUTDOWN = new ProxyStartAndShutdown(); private static class ProxyStartAndShutdown extends AbstractStartAndShutdown { - @Override public void appendStartAndShutdown(StartAndShutdown startAndShutdown) { + @Override + public void appendStartAndShutdown(StartAndShutdown startAndShutdown) { super.appendStartAndShutdown(startAndShutdown); } } @@ -59,30 +60,29 @@ public class ProxyStartup { // create and start grpcServer GrpcServer grpcServer = createGrpcServer(); - proxyStartAndShutdown.appendStartAndShutdown(grpcServer); + PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(grpcServer); // health check server final HealthCheckServer healthCheckServer = new HealthCheckServer(); - proxyStartAndShutdown.appendStartAndShutdown(healthCheckServer); + PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(healthCheckServer); Runtime.getRuntime().addShutdownHook(new Thread(() -> { - log.info("try to shutdown server"); - + LOGGER.info("try to shutdown server"); try { - proxyStartAndShutdown.shutdown(); + PROXY_START_AND_SHUTDOWN.shutdown(); } catch (Exception e) { - log.error("err when shutdown proxy", e); + LOGGER.error("err when shutdown proxy", e); } })); } catch (Exception e) { System.err.println("find a unexpect err." + e); e.printStackTrace(); - log.error("find a unexpect err.", e); + LOGGER.error("find a unexpect err.", e); System.exit(1); } System.out.printf("%s%n", new Date() + " rmq-proxy startup successfully"); - log.info(new Date() + "rmq-proxy startup successfully"); + LOGGER.info(new Date() + "rmq-proxy startup successfully"); } private static GrpcServer createGrpcServer() throws Exception { @@ -93,15 +93,17 @@ public class ProxyStartup { } else if (ProxyMode.isLocalMode(proxyModeStr)) { BrokerController brokerController = createBrokerController(); StartAndShutdown brokerControllerWrapper = new StartAndShutdown() { - @Override public void start() throws Exception { + @Override + public void start() throws Exception { brokerController.start(); } - @Override public void shutdown() throws Exception { + @Override + public void shutdown() throws Exception { brokerController.shutdown(); } }; - proxyStartAndShutdown.appendStartAndShutdown(brokerControllerWrapper); + PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(brokerControllerWrapper); grpcService = new LocalGrpcService(brokerController); } else { throw new IllegalArgumentException("try to start grpc server with wrong mode, use 'local' or 'cluster'"); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java index 8b0739efd3..bdf2ef601f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/InvocationContext.java @@ -21,9 +21,9 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; public class InvocationContext { - final private R request; - final private CompletableFuture response; - final private long timestamp = System.currentTimeMillis(); + private final R request; + private final CompletableFuture response; + private final long timestamp = System.currentTimeMillis(); public InvocationContext(R req, CompletableFuture resp) { request = req; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java index 99ab1da516..65bd2ad5eb 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java @@ -52,11 +52,17 @@ public class GrpcClientChannel extends SimpleChannel { this.pollCommandResponseFutureRef.set(future); } - public static GrpcClientChannel create(ChannelManager channelManager, String group, String clientId, PollCommandResponseManager manager) { + public static GrpcClientChannel create( + ChannelManager channelManager, + String group, + String clientId, + PollCommandResponseManager manager + ) { GrpcClientChannel channel = channelManager.createChannel( buildKey(group, clientId), () -> new GrpcClientChannel(group, clientId, manager), - GrpcClientChannel.class); + GrpcClientChannel.class + ); channelManager.addGroupClientId(group, clientId); return channel; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java index d106f8f0d2..1cbb003610 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java @@ -46,9 +46,11 @@ public class HeaderInterceptor implements ServerInterceptor { private String parseSocketAddress(SocketAddress socketAddress) { if (socketAddress instanceof InetSocketAddress) { InetSocketAddress inetSocketAddress = (InetSocketAddress) socketAddress; - return HostAndPort.fromParts(inetSocketAddress.getAddress() - .getHostAddress(), inetSocketAddress.getPort()) - .toString(); + return HostAndPort.fromParts( + inetSocketAddress.getAddress() + .getHostAddress(), + inetSocketAddress.getPort() + ).toString(); } return ""; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 889da0e2b4..1846efe19d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -181,8 +181,9 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc return this.clientService.pollCommand(ctx, request); } - @Override public CompletableFuture reportThreadStackTrace(Context ctx, - ReportThreadStackTraceRequest request) { + @Override + public CompletableFuture reportThreadStackTrace(Context ctx, + ReportThreadStackTraceRequest request) { return null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index ce4e7f78b1..d4c2cce67b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -136,7 +136,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo this.appendStartAndShutdown(new LocalGrpcServiceStartAndShutdown()); } - @Override public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { + @Override + public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { return this.routeService.queryRoute(ctx, request); } @@ -257,7 +258,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return future; } - @Override public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { + @Override + public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); AckMessageRequestHeader requestHeader = Converter.buildAckMessageRequestHeader(request); @@ -282,7 +284,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return future; } - @Override public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { + @Override + public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); @@ -364,7 +367,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return future; } - @Override public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) { + @Override + public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) { Partition partition = request.getPartition(); String topicName = Converter.getResourceNameWithNamespace(partition.getTopic()); int queueId = partition.getId(); @@ -386,7 +390,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo .build()); } - @Override public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) { + @Override + public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) { long timeRemaining = Context.current() .getDeadline() .timeRemaining(TimeUnit.MILLISECONDS); @@ -419,7 +424,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return future; } - @Override public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { + @Override + public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { String clientId = request.getClientId(); CompletableFuture future = new CompletableFuture<>(); switch (request.getGroupCase()) { @@ -479,6 +485,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo @Override public CompletableFuture reportMessageConsumptionResult(Context ctx, ReportMessageConsumptionResultRequest request) { + String commandId = request.getCommandId(); PollCommandResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId); if (pollCommandResponseFuture != null) { @@ -497,7 +504,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo .build()); } - @Override public CompletableFuture notifyClientTermination(Context ctx, + @Override + public CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request) { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel); @@ -512,7 +520,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return new CompletableFuture<>(); } - @Override public CompletableFuture changeInvisibleDuration(Context ctx, + @Override + public CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request) { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); 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 9872bd35c3..69aa281186 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 @@ -74,7 +74,6 @@ 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); }