From 76cf9dbdc18cb65eaa01be74a96cb1a0673b4764 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Wed, 23 Mar 2022 11:40:35 +0800 Subject: [PATCH] [ISSUE #3949] For passing check style. --- ...entAPIExtImpl.java => MQClientAPIExt.java} | 125 +++++++++++------- .../connector/AbstractForwardClient.java | 13 +- .../proxy/connector/DefaultForwardClient.java | 4 +- .../proxy/connector/ForwardProducer.java | 4 +- .../proxy/connector/ForwardReadConsumer.java | 4 +- .../proxy/connector/ForwardWriteConsumer.java | 4 +- .../factory/AbstractMQClientFactory.java | 16 +-- .../factory/ForwardClientFactory.java | 6 +- .../AuthenticationInterceptor.java | 8 +- .../grpc/service/ClusterGrpcService.java | 9 +- 10 files changed, 115 insertions(+), 78 deletions(-) rename client/src/main/java/org/apache/rocketmq/client/impl/{MQClientAPIExtImpl.java => MQClientAPIExt.java} (85%) diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExt.java similarity index 85% rename from client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java rename to client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExt.java index a9474fd37e..c440cd6c0b 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExt.java @@ -57,16 +57,18 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -public class MQClientAPIExtImpl { +public class MQClientAPIExt { + private static final Logger LOGGER = LoggerFactory.getLogger(MQClientAPIExt.class); - private static final Logger log = LoggerFactory.getLogger(MQClientAPIExtImpl.class); - - private final MQClientAPIImpl mqClientAPI; private final ClientConfig clientConfig; + private final MQClientAPIImpl mqClientAPI; - public MQClientAPIExtImpl(NettyClientConfig nettyClientConfig, + public MQClientAPIExt( + ClientConfig clientConfig, + NettyClientConfig nettyClientConfig, ClientRemotingProcessor clientRemotingProcessor, - RPCHook rpcHook, ClientConfig clientConfig) { + RPCHook rpcHook + ) { this.clientConfig = clientConfig; this.mqClientAPI = new MQClientAPIImpl(nettyClientConfig, clientRemotingProcessor, rpcHook, clientConfig); } @@ -86,7 +88,7 @@ public class MQClientAPIExtImpl { public boolean updateNameServerAddressList() { if (this.clientConfig.getNamesrvAddr() != null) { this.mqClientAPI.updateNameServerAddressList(this.clientConfig.getNamesrvAddr()); - log.info("user specified name server address: {}", this.clientConfig.getNamesrvAddr()); + LOGGER.info("user specified name server address: {}", this.clientConfig.getNamesrvAddr()); return true; } return false; @@ -109,8 +111,11 @@ public class MQClientAPIExtImpl { return this.mqClientAPI.getRemotingClient(); } - public CompletableFuture sendHeartbeat(String brokerAddr, HeartbeatData heartbeatData, - long timeoutMillis) { + public CompletableFuture sendHeartbeat( + String brokerAddr, + HeartbeatData heartbeatData, + long timeoutMillis + ) { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.HEART_BEAT, null); request.setLanguage(clientConfig.getLanguage()); request.setBody(heartbeatData.encode()); @@ -140,9 +145,13 @@ public class MQClientAPIExtImpl { this.mqClientAPI.endTransactionOneway(brokerAddr, requestHeader, remark, timeoutMillis); } - public CompletableFuture sendMessage(String brokerAddr, String brokerName, Message msg, - SendMessageRequestHeader requestHeader, long timeoutMillis) { - + public CompletableFuture sendMessage( + String brokerAddr, + String brokerName, + Message msg, + SendMessageRequestHeader requestHeader, + long timeoutMillis + ) { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.SEND_MESSAGE, requestHeader); request.setBody(msg.getBody()); @@ -166,9 +175,11 @@ public class MQClientAPIExtImpl { return future; } - public CompletableFuture sendMessageBack(String brokerAddr, + public CompletableFuture sendMessageBack( + String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, - long timeoutMillis) { + long timeoutMillis + ) { CompletableFuture future = new CompletableFuture<>(); try { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader); @@ -186,9 +197,12 @@ public class MQClientAPIExtImpl { return future; } - public CompletableFuture popMessage(String brokerAddr, String brokerName, + public CompletableFuture popMessage( + String brokerAddr, + String brokerName, PopMessageRequestHeader requestHeader, - long timeoutMillis) { + long timeoutMillis + ) { CompletableFuture future = new CompletableFuture<>(); try { this.mqClientAPI.popMessageAsync(brokerName, brokerAddr, requestHeader, timeoutMillis, new PopCallback() { @@ -208,8 +222,11 @@ public class MQClientAPIExtImpl { return future; } - public CompletableFuture ackMessage(String brokerAddr, AckMessageRequestHeader requestHeader, - long timeoutMillis) { + public CompletableFuture ackMessage( + String brokerAddr, + AckMessageRequestHeader requestHeader, + long timeoutMillis + ) { CompletableFuture future = new CompletableFuture<>(); try { this.mqClientAPI.ackMessageAsync(brokerAddr, timeoutMillis, new AckCallback() { @@ -229,56 +246,73 @@ public class MQClientAPIExtImpl { return future; } - public CompletableFuture changeInvisibleTimeAsync(String brokerAddr, String brokerName, - ChangeInvisibleTimeRequestHeader requestHeader, long timeoutMillis) { + public CompletableFuture changeInvisibleTimeAsync( + String brokerAddr, + String brokerName, + ChangeInvisibleTimeRequestHeader requestHeader, + long timeoutMillis + ) { CompletableFuture future = new CompletableFuture<>(); try { - this.mqClientAPI.changeInvisibleTimeAsync(brokerName, brokerAddr, requestHeader, timeoutMillis, new AckCallback() { - @Override - public void onSuccess(AckResult ackResult) { - future.complete(ackResult); - } + this.mqClientAPI.changeInvisibleTimeAsync(brokerName, brokerAddr, requestHeader, timeoutMillis, + new AckCallback() { + @Override + public void onSuccess(AckResult ackResult) { + future.complete(ackResult); + } - @Override - public void onException(Throwable t) { - future.completeExceptionally(t); + @Override + public void onException(Throwable t) { + future.completeExceptionally(t); + } } - }); + ); } catch (Throwable t) { future.completeExceptionally(t); } return future; } - public CompletableFuture pullMessage(String brokerAddr, PullMessageRequestHeader requestHeader, - long timeoutMillis) { + public CompletableFuture pullMessage( + String brokerAddr, + PullMessageRequestHeader requestHeader, + long timeoutMillis + ) { CompletableFuture future = new CompletableFuture<>(); try { - this.mqClientAPI.pullMessage(brokerAddr, requestHeader, timeoutMillis, CommunicationMode.ASYNC, new PullCallback() { - @Override - public void onSuccess(PullResult pullResult) { - future.complete(pullResult); - } + this.mqClientAPI.pullMessage(brokerAddr, requestHeader, timeoutMillis, CommunicationMode.ASYNC, + new PullCallback() { + @Override + public void onSuccess(PullResult pullResult) { + future.complete(pullResult); + } - @Override - public void onException(Throwable t) { - future.completeExceptionally(t); + @Override + public void onException(Throwable t) { + future.completeExceptionally(t); + } } - }); + ); } catch (Throwable t) { future.completeExceptionally(t); } return future; } - public void updateConsumerOffsetOneWay(String brokerAddr, UpdateConsumerOffsetRequestHeader header, - long timeoutMillis) throws InterruptedException, RemotingException { + public void updateConsumerOffsetOneWay( + String brokerAddr, + UpdateConsumerOffsetRequestHeader header, + long timeoutMillis + ) throws InterruptedException, RemotingException { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.UPDATE_CONSUMER_OFFSET, header); this.getRemotingClient().invokeOneway(brokerAddr, request, timeoutMillis); } - public CompletableFuture> getConsumerListByGroup(String brokerAddr, GetConsumerListByGroupRequestHeader requestHeader, - long timeoutMillis) { + public CompletableFuture> getConsumerListByGroup( + String brokerAddr, + GetConsumerListByGroupRequestHeader requestHeader, + long timeoutMillis + ) { RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_CONSUMER_LIST_BY_GROUP, requestHeader); CompletableFuture> future = new CompletableFuture<>(); @@ -317,7 +351,8 @@ public class MQClientAPIExtImpl { return future; } - public TopicRouteData getTopicRouteInfoFromNameServer(String topic, long timeoutMillis) throws RemotingException, InterruptedException, MQClientException { + public TopicRouteData getTopicRouteInfoFromNameServer(String topic, long timeoutMillis) + throws RemotingException, InterruptedException, MQClientException { return this.mqClientAPI.getTopicRouteInfoFromNameServer(topic, timeoutMillis); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java index 0fa3c3746b..2f08e7e67c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java @@ -17,14 +17,14 @@ package org.apache.rocketmq.proxy.connector; import java.util.concurrent.ThreadLocalRandom; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; -import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.proxy.common.StartAndShutdown; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; public abstract class AbstractForwardClient implements StartAndShutdown { private final ForwardClientFactory forwardClientFactory; - private MQClientAPIExtImpl[] clients; + private MQClientAPIExt[] clients; public AbstractForwardClient(ForwardClientFactory forwardClientFactory) { this.forwardClientFactory = forwardClientFactory; @@ -32,11 +32,11 @@ public abstract class AbstractForwardClient implements StartAndShutdown { protected abstract int getClientNum(); - protected abstract MQClientAPIExtImpl createNewClient(ForwardClientFactory forwardClientFactory, String name); + protected abstract MQClientAPIExt createNewClient(ForwardClientFactory forwardClientFactory, String name); protected abstract String getNamePrefix(); - protected MQClientAPIExtImpl getClient() { + protected MQClientAPIExt getClient() { if (clients.length == 1) { return this.clients[0]; } @@ -46,7 +46,8 @@ public abstract class AbstractForwardClient implements StartAndShutdown { @Override public void start() throws Exception { int clientCount = getClientNum(); - this.clients = new MQClientAPIExtImpl[clientCount]; + this.clients = new MQClientAPIExt[clientCount]; + for (int i = 0; i < clientCount; i++) { String name = getNamePrefix() + "N_" + i; clients[i] = createNewClient(forwardClientFactory, name); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java index b1e032c2f4..81d180e66e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java @@ -19,7 +19,7 @@ package org.apache.rocketmq.proxy.connector; import java.util.List; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.exception.MQClientException; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.common.protocol.header.GetConsumerListByGroupRequestHeader; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.proxy.config.ConfigurationManager; @@ -39,7 +39,7 @@ public class DefaultForwardClient extends AbstractForwardClient { } @Override - protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getDefaultForwardClientWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java index 8facf9a8e4..2f6d504c72 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java @@ -17,7 +17,7 @@ package org.apache.rocketmq.proxy.connector; import java.util.concurrent.CompletableFuture; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.message.Message; @@ -46,7 +46,7 @@ public class ForwardProducer extends AbstractForwardClient { } @Override - protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { double sendClientWorkerFactor = ConfigurationManager.getProxyConfig().getForwardProducerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * sendClientWorkerFactor); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java index 0aa0b30d70..a8faa303b3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java @@ -19,7 +19,7 @@ package org.apache.rocketmq.proxy.connector; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.PopResult; import org.apache.rocketmq.client.consumer.PullResult; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader; import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; @@ -39,7 +39,7 @@ public class ForwardReadConsumer extends AbstractForwardClient { } @Override - protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getForwardConsumerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java index 5d3f5c2035..fc67796796 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java @@ -18,7 +18,7 @@ package org.apache.rocketmq.proxy.connector; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.AckResult; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader; import org.apache.rocketmq.common.protocol.header.UpdateConsumerOffsetRequestHeader; @@ -40,7 +40,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient { } @Override - protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getForwardConsumerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java index c7417d4d46..6f126e9200 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java @@ -20,10 +20,10 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.client.ClientConfig; import org.apache.rocketmq.client.impl.ClientRemotingProcessor; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.remoting.RPCHook; -public abstract class AbstractMQClientFactory extends AbstractClientFactory { +public abstract class AbstractMQClientFactory extends AbstractClientFactory { public AbstractMQClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) { @@ -33,20 +33,20 @@ public abstract class AbstractMQClientFactory extends AbstractClientFactory ServerCall.Listener interceptCall(ServerCall call, Metadata headers, - ServerCallHandler next) { - return new ForwardingServerCallListener.SimpleForwardingServerCallListener(next.startCall(call, headers)) { + public ServerCall.Listener interceptCall(ServerCall call, Metadata headers, + ServerCallHandler next) { + return new ForwardingServerCallListener.SimpleForwardingServerCallListener(next.startCall(call, headers)) { @Override - public void onMessage(ReqT message) { + public void onMessage(R message) { GeneratedMessageV3 messageV3 = (GeneratedMessageV3) message; MetadataHeader metadataHeader = MetadataHeader.builder() .remoteAddress(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.REMOTE_ADDRESS)) 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 1846efe19d..63bee24262 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 @@ -124,10 +124,11 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc @Override public CompletableFuture healthCheck(Context ctx, HealthCheckRequest request) { - final HealthCheckResponse response = HealthCheckResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) - .build(); - return CompletableFuture.completedFuture(response); + return CompletableFuture.completedFuture( + HealthCheckResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) + .build() + ); } @Override