[ISSUE #3949] For passing check style.

This commit is contained in:
Jixiang.jjx
2022-07-13 11:29:13 +08:00
committed by zhouxiang
parent 15d181cd70
commit 76cf9dbdc1
10 changed files with 115 additions and 78 deletions
@@ -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<Integer> sendHeartbeat(String brokerAddr, HeartbeatData heartbeatData,
long timeoutMillis) {
public CompletableFuture<Integer> 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<SendResult> sendMessage(String brokerAddr, String brokerName, Message msg,
SendMessageRequestHeader requestHeader, long timeoutMillis) {
public CompletableFuture<SendResult> 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<RemotingCommand> sendMessageBack(String brokerAddr,
public CompletableFuture<RemotingCommand> sendMessageBack(
String brokerAddr,
ConsumerSendMsgBackRequestHeader requestHeader,
long timeoutMillis) {
long timeoutMillis
) {
CompletableFuture<RemotingCommand> 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<PopResult> popMessage(String brokerAddr, String brokerName,
public CompletableFuture<PopResult> popMessage(
String brokerAddr,
String brokerName,
PopMessageRequestHeader requestHeader,
long timeoutMillis) {
long timeoutMillis
) {
CompletableFuture<PopResult> future = new CompletableFuture<>();
try {
this.mqClientAPI.popMessageAsync(brokerName, brokerAddr, requestHeader, timeoutMillis, new PopCallback() {
@@ -208,8 +222,11 @@ public class MQClientAPIExtImpl {
return future;
}
public CompletableFuture<AckResult> ackMessage(String brokerAddr, AckMessageRequestHeader requestHeader,
long timeoutMillis) {
public CompletableFuture<AckResult> ackMessage(
String brokerAddr,
AckMessageRequestHeader requestHeader,
long timeoutMillis
) {
CompletableFuture<AckResult> future = new CompletableFuture<>();
try {
this.mqClientAPI.ackMessageAsync(brokerAddr, timeoutMillis, new AckCallback() {
@@ -229,56 +246,73 @@ public class MQClientAPIExtImpl {
return future;
}
public CompletableFuture<AckResult> changeInvisibleTimeAsync(String brokerAddr, String brokerName,
ChangeInvisibleTimeRequestHeader requestHeader, long timeoutMillis) {
public CompletableFuture<AckResult> changeInvisibleTimeAsync(
String brokerAddr,
String brokerName,
ChangeInvisibleTimeRequestHeader requestHeader,
long timeoutMillis
) {
CompletableFuture<AckResult> 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<PullResult> pullMessage(String brokerAddr, PullMessageRequestHeader requestHeader,
long timeoutMillis) {
public CompletableFuture<PullResult> pullMessage(
String brokerAddr,
PullMessageRequestHeader requestHeader,
long timeoutMillis
) {
CompletableFuture<PullResult> 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<List<String>> getConsumerListByGroup(String brokerAddr, GetConsumerListByGroupRequestHeader requestHeader,
long timeoutMillis) {
public CompletableFuture<List<String>> getConsumerListByGroup(
String brokerAddr,
GetConsumerListByGroupRequestHeader requestHeader,
long timeoutMillis
) {
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_CONSUMER_LIST_BY_GROUP, requestHeader);
CompletableFuture<List<String>> 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);
}
@@ -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);
@@ -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);
@@ -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);
@@ -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);
@@ -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);
@@ -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<MQClientAPIExtImpl> {
public abstract class AbstractMQClientFactory extends AbstractClientFactory<MQClientAPIExt> {
public AbstractMQClientFactory(ScheduledExecutorService scheduledExecutorService,
RPCHook rpcHook) {
@@ -33,20 +33,20 @@ public abstract class AbstractMQClientFactory extends AbstractClientFactory<MQCl
protected abstract ClientRemotingProcessor createClientRemotingProcessor();
@Override
protected MQClientAPIExtImpl newOne(String instanceName, RPCHook rpcHook, int bootstrapWorkerThreads) {
protected MQClientAPIExt newOne(String instanceName, RPCHook rpcHook, int bootstrapWorkerThreads) {
ClientConfig clientConfig = new ClientConfig();
clientConfig.setInstanceName(instanceName);
return new MQClientAPIExtImpl(
return new MQClientAPIExt(
clientConfig,
createNettyClientConfig(bootstrapWorkerThreads),
createClientRemotingProcessor(),
rpcHook,
clientConfig
rpcHook
);
}
@Override
protected boolean tryStart(MQClientAPIExtImpl client) {
protected boolean tryStart(MQClientAPIExt client) {
if (!client.updateNameServerAddressList()) {
this.scheduledExecutorService.scheduleAtFixedRate(
client::fetchNameServerAddr,
@@ -60,7 +60,7 @@ public abstract class AbstractMQClientFactory extends AbstractClientFactory<MQCl
}
@Override
protected void shutdown(MQClientAPIExtImpl client) {
protected void shutdown(MQClientAPIExt client) {
client.shutdown();
}
}
@@ -21,7 +21,7 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.client.ClientConfig;
import org.apache.rocketmq.client.impl.MQClientAPIExtImpl;
import org.apache.rocketmq.client.impl.MQClientAPIExt;
import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker;
import org.apache.rocketmq.proxy.common.StartAndShutdown;
@@ -56,11 +56,11 @@ public class ForwardClientFactory implements StartAndShutdown {
}
}
public MQClientAPIExtImpl getMQClient(String instanceName, int bootstrapWorkerThreads) {
public MQClientAPIExt getMQClient(String instanceName, int bootstrapWorkerThreads) {
return mqClientFactory.getOne(instanceName, bootstrapWorkerThreads);
}
public MQClientAPIExtImpl getTransactionalProducer(String instanceName, int bootstrapWorkerThreads) {
public MQClientAPIExt getTransactionalProducer(String instanceName, int bootstrapWorkerThreads) {
return transactionalProducerFactory.getOne(instanceName, bootstrapWorkerThreads);
}
@@ -37,11 +37,11 @@ public class AuthenticationInterceptor implements ServerInterceptor {
}
@Override
public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> call, Metadata headers,
ServerCallHandler<ReqT, RespT> next) {
return new ForwardingServerCallListener.SimpleForwardingServerCallListener<ReqT>(next.startCall(call, headers)) {
public <R, W> ServerCall.Listener<R> interceptCall(ServerCall<R, W> call, Metadata headers,
ServerCallHandler<R, W> next) {
return new ForwardingServerCallListener.SimpleForwardingServerCallListener<R>(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))
@@ -124,10 +124,11 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
@Override
public CompletableFuture<HealthCheckResponse> 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